Optimize scheduler for minimal SQLite database locks
- Redesigned scheduler to use separate short-lived sessions for lock operations - Lock acquisition and release now use independent database connections - Each scheduled action processed in its own transaction to minimize lock time - Error handling isolates failed actions to prevent blocking others - Re-enabled scheduler with optimized locking strategy Key improvements: - acquire_scheduler_lock() uses own connection, disposes immediately - release_scheduler_lock() uses own connection, disposes immediately - process_single_scheduled_action() uses transaction-per-action approach - Prevents long-running transactions that were blocking web requests russell@unturf.com is the boss
This commit is contained in:
parent
3cede5ba07
commit
6a3dd08728
2 changed files with 120 additions and 77 deletions
|
|
@ -63,7 +63,7 @@ deploy:
|
||||||
|
|
||||||
[Service]
|
[Service]
|
||||||
WorkingDirectory=${APP_DIR}
|
WorkingDirectory=${APP_DIR}
|
||||||
Environment="PATH=${VENV_DIR}/bin" "DISABLE_SCHEDULER=true"
|
Environment="PATH=${VENV_DIR}/bin"
|
||||||
ExecStart=${VENV_DIR}/bin/uwsgi \\
|
ExecStart=${VENV_DIR}/bin/uwsgi \\
|
||||||
--master \\
|
--master \\
|
||||||
--enable-threads \\
|
--enable-threads \\
|
||||||
|
|
|
||||||
195
app.py
195
app.py
|
|
@ -622,49 +622,72 @@ def reader_required(view_func):
|
||||||
################################################################################
|
################################################################################
|
||||||
|
|
||||||
|
|
||||||
def acquire_scheduler_lock(dbsession):
|
def acquire_scheduler_lock():
|
||||||
"""Acquire the scheduler lock."""
|
"""Acquire the scheduler lock with minimal transaction time."""
|
||||||
process_id = f"{os.getpid()}_{int(time.time())}"
|
process_id = f"{os.getpid()}_{int(time.time())}"
|
||||||
|
|
||||||
|
# Use a separate short-lived session just for the lock
|
||||||
|
main_engine = create_engine(
|
||||||
|
DB_URL,
|
||||||
|
connect_args={"check_same_thread": False},
|
||||||
|
poolclass=StaticPool,
|
||||||
|
)
|
||||||
|
MainSessionFactory = sessionmaker(bind=main_engine)
|
||||||
|
|
||||||
|
try:
|
||||||
|
with MainSessionFactory() as dbsession:
|
||||||
|
# Try to get existing lock
|
||||||
|
lock = dbsession.query(SchedulerLock).first()
|
||||||
|
|
||||||
# Try to get existing lock
|
if not lock:
|
||||||
lock = dbsession.query(SchedulerLock).first()
|
# Create new lock
|
||||||
|
lock = SchedulerLock(
|
||||||
|
locked=True, locked_at=datetime.datetime.utcnow(), process_id=process_id
|
||||||
|
)
|
||||||
|
dbsession.add(lock)
|
||||||
|
dbsession.commit()
|
||||||
|
return True
|
||||||
|
|
||||||
if not lock:
|
# Check if lock is stale (older than 5 minutes)
|
||||||
# Create new lock
|
if lock.locked and lock.locked_at:
|
||||||
lock = SchedulerLock(
|
if datetime.datetime.utcnow() - lock.locked_at > datetime.timedelta(minutes=5):
|
||||||
locked=True, locked_at=datetime.datetime.utcnow(), process_id=process_id
|
log.warning("Detected stale scheduler lock, releasing it")
|
||||||
)
|
lock.locked = False
|
||||||
dbsession.add(lock)
|
lock.locked_at = None
|
||||||
dbsession.commit()
|
lock.process_id = None
|
||||||
return True
|
|
||||||
|
|
||||||
# Check if lock is stale (older than 5 minutes)
|
if not lock.locked:
|
||||||
if lock.locked and lock.locked_at:
|
lock.locked = True
|
||||||
if datetime.datetime.utcnow() - lock.locked_at > datetime.timedelta(minutes=5):
|
lock.locked_at = datetime.datetime.utcnow()
|
||||||
log.warning("Detected stale scheduler lock, releasing it")
|
lock.process_id = process_id
|
||||||
lock.locked = False
|
dbsession.commit()
|
||||||
lock.locked_at = None
|
return True
|
||||||
lock.process_id = None
|
|
||||||
dbsession.commit()
|
|
||||||
|
|
||||||
if not lock.locked:
|
return False
|
||||||
lock.locked = True
|
finally:
|
||||||
lock.locked_at = datetime.datetime.utcnow()
|
main_engine.dispose()
|
||||||
lock.process_id = process_id
|
|
||||||
dbsession.commit()
|
|
||||||
return True
|
|
||||||
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
def release_scheduler_lock(dbsession):
|
def release_scheduler_lock():
|
||||||
"""Release the scheduler lock."""
|
"""Release the scheduler lock with minimal transaction time."""
|
||||||
lock = dbsession.query(SchedulerLock).first()
|
# Use a separate short-lived session just for the lock
|
||||||
if lock and lock.locked:
|
main_engine = create_engine(
|
||||||
lock.locked = False
|
DB_URL,
|
||||||
lock.locked_at = None
|
connect_args={"check_same_thread": False},
|
||||||
lock.process_id = None
|
poolclass=StaticPool,
|
||||||
dbsession.commit()
|
)
|
||||||
|
MainSessionFactory = sessionmaker(bind=main_engine)
|
||||||
|
|
||||||
|
try:
|
||||||
|
with MainSessionFactory() as dbsession:
|
||||||
|
lock = dbsession.query(SchedulerLock).first()
|
||||||
|
if lock and lock.locked:
|
||||||
|
lock.locked = False
|
||||||
|
lock.locked_at = None
|
||||||
|
lock.process_id = None
|
||||||
|
dbsession.commit()
|
||||||
|
finally:
|
||||||
|
main_engine.dispose()
|
||||||
|
|
||||||
|
|
||||||
def process_scheduled_actions():
|
def process_scheduled_actions():
|
||||||
|
|
@ -678,23 +701,27 @@ def process_scheduled_actions():
|
||||||
main_dbsession = MainSessionFactory()
|
main_dbsession = MainSessionFactory()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not acquire_scheduler_lock(main_dbsession):
|
if not acquire_scheduler_lock():
|
||||||
return # Another process is already running
|
return # Another process is already running
|
||||||
|
|
||||||
log.info("Scheduler acquired lock, processing scheduled actions...")
|
log.info("Scheduler acquired lock, processing scheduled actions...")
|
||||||
|
|
||||||
# Get all namespaces
|
# Get all namespaces with the existing session
|
||||||
namespaces = main_dbsession.query(Namespace).all()
|
namespaces = main_dbsession.query(Namespace).all()
|
||||||
|
|
||||||
|
# Process each namespace separately to avoid long transactions
|
||||||
for namespace in namespaces:
|
for namespace in namespaces:
|
||||||
process_namespace_scheduled_actions(namespace.id)
|
try:
|
||||||
|
process_namespace_scheduled_actions(namespace.id)
|
||||||
|
except Exception as e:
|
||||||
|
log.error(f"Error processing namespace {namespace.id}: {e}")
|
||||||
|
|
||||||
log.info("Finished processing scheduled actions")
|
log.info("Finished processing scheduled actions")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
log.error(f"Error in scheduler: {e}")
|
log.error(f"Error in scheduler: {e}")
|
||||||
finally:
|
finally:
|
||||||
release_scheduler_lock(main_dbsession)
|
release_scheduler_lock()
|
||||||
main_dbsession.close()
|
main_dbsession.close()
|
||||||
main_engine.dispose()
|
main_engine.dispose()
|
||||||
|
|
||||||
|
|
@ -731,43 +758,9 @@ def process_namespace_scheduled_actions(namespace_id):
|
||||||
.all()
|
.all()
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Process each action individually with short transactions
|
||||||
for action in pending_actions:
|
for action in pending_actions:
|
||||||
try:
|
process_single_scheduled_action(namespace_dbsession, action)
|
||||||
media = (
|
|
||||||
namespace_dbsession.query(Media)
|
|
||||||
.filter_by(id=action.media_id)
|
|
||||||
.first()
|
|
||||||
)
|
|
||||||
if not media:
|
|
||||||
action.status = "failed"
|
|
||||||
continue
|
|
||||||
|
|
||||||
if action.action_type == "delete":
|
|
||||||
namespace_dbsession.delete(media)
|
|
||||||
log.info(f"Scheduled delete executed for media {media.filename}")
|
|
||||||
elif action.action_type == "set_public":
|
|
||||||
media.visibility = "public"
|
|
||||||
log.info(
|
|
||||||
f"Scheduled set_public executed for media {media.filename}"
|
|
||||||
)
|
|
||||||
elif action.action_type == "set_private":
|
|
||||||
media.visibility = "private"
|
|
||||||
log.info(
|
|
||||||
f"Scheduled set_private executed for media {media.filename}"
|
|
||||||
)
|
|
||||||
elif action.action_type == "set_unlisted":
|
|
||||||
media.visibility = "unlisted"
|
|
||||||
log.info(
|
|
||||||
f"Scheduled set_unlisted executed for media {media.filename}"
|
|
||||||
)
|
|
||||||
|
|
||||||
action.status = "completed"
|
|
||||||
|
|
||||||
except Exception as e:
|
|
||||||
log.error(f"Error processing scheduled action {action.id}: {e}")
|
|
||||||
action.status = "failed"
|
|
||||||
|
|
||||||
namespace_dbsession.commit()
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
log.error(f"Error processing namespace {namespace_id}: {e}")
|
log.error(f"Error processing namespace {namespace_id}: {e}")
|
||||||
|
|
@ -778,6 +771,56 @@ def process_namespace_scheduled_actions(namespace_id):
|
||||||
namespace_engine.dispose()
|
namespace_engine.dispose()
|
||||||
|
|
||||||
|
|
||||||
|
def process_single_scheduled_action(dbsession, action):
|
||||||
|
"""Process a single scheduled action with minimal transaction time."""
|
||||||
|
try:
|
||||||
|
# Start a new transaction for this action
|
||||||
|
with dbsession.begin():
|
||||||
|
# Refresh the action to get latest state
|
||||||
|
dbsession.refresh(action)
|
||||||
|
|
||||||
|
# Skip if already processed
|
||||||
|
if action.status != "pending":
|
||||||
|
return
|
||||||
|
|
||||||
|
media = (
|
||||||
|
dbsession.query(Media)
|
||||||
|
.filter_by(id=action.media_id)
|
||||||
|
.first()
|
||||||
|
)
|
||||||
|
if not media:
|
||||||
|
action.status = "failed"
|
||||||
|
action.completed_at = datetime.datetime.utcnow()
|
||||||
|
return
|
||||||
|
|
||||||
|
if action.action_type == "delete":
|
||||||
|
dbsession.delete(media)
|
||||||
|
log.info(f"Scheduled delete executed for media {media.filename}")
|
||||||
|
elif action.action_type == "set_public":
|
||||||
|
media.visibility = "public"
|
||||||
|
log.info(f"Scheduled set_public executed for media {media.filename}")
|
||||||
|
elif action.action_type == "set_private":
|
||||||
|
media.visibility = "private"
|
||||||
|
log.info(f"Scheduled set_private executed for media {media.filename}")
|
||||||
|
elif action.action_type == "set_unlisted":
|
||||||
|
media.visibility = "unlisted"
|
||||||
|
log.info(f"Scheduled set_unlisted executed for media {media.filename}")
|
||||||
|
|
||||||
|
action.status = "completed"
|
||||||
|
action.completed_at = datetime.datetime.utcnow()
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
log.error(f"Error processing scheduled action {action.id}: {e}")
|
||||||
|
try:
|
||||||
|
# Try to mark as failed in a separate transaction
|
||||||
|
with dbsession.begin():
|
||||||
|
dbsession.refresh(action)
|
||||||
|
action.status = "failed"
|
||||||
|
action.completed_at = datetime.datetime.utcnow()
|
||||||
|
except Exception as inner_e:
|
||||||
|
log.error(f"Failed to mark action {action.id} as failed: {inner_e}")
|
||||||
|
|
||||||
|
|
||||||
# Global variables for scheduler thread
|
# Global variables for scheduler thread
|
||||||
scheduler_thread = None
|
scheduler_thread = None
|
||||||
scheduler_running = False
|
scheduler_running = False
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue