From d05279b4107d01b563b2876bda1b4eee560df620 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sun, 23 Aug 2026 12:21:45 +0800 Subject: [PATCH 1/4] feat(cleanup): pending task auto cleanup via Celery Beat - generation_task_repository: add cleanup_stale_pending() method marks status=pending + created_at > timeout as failed with error_message="pending timeout: auto cleanup" - _startup.py: add cleanup_stale_pending_tasks() function with PENDING_TASK_TIMEOUT_MINUTES = 30 constant - cleanup.py: new Celery scheduled task worker.cleanup_stale_pending_tasks - celery_app.py: add beat_schedule, runs every 10 minutes (600s) - entrypoint-worker.sh: add -B flag for embedded celery beat - 6 unit tests for cleanup_stale_pending --- apps/worker/worker_app/celery_app.py | 10 ++ apps/worker/worker_app/tasks/_startup.py | 42 ++++- apps/worker/worker_app/tasks/cleanup.py | 36 +++++ infra/docker/entrypoint-worker.sh | 3 + .../generation_task_repository.py | 37 +++++ .../test_generation_task_pending_cleanup.py | 150 ++++++++++++++++++ 6 files changed, 275 insertions(+), 3 deletions(-) create mode 100644 apps/worker/worker_app/tasks/cleanup.py create mode 100644 tests/unit/test_generation_task_pending_cleanup.py diff --git a/apps/worker/worker_app/celery_app.py b/apps/worker/worker_app/celery_app.py index 6e51e61ff..2cf626277 100755 --- a/apps/worker/worker_app/celery_app.py +++ b/apps/worker/worker_app/celery_app.py @@ -19,4 +19,14 @@ celery_app.conf.imports = ( "worker_app.tasks.batch_download", "worker_app.tasks._startup", "apps.worker.video_processing.dedup", + "worker_app.tasks.cleanup", ) + +# Celery Beat 定时任务调度 +celery_app.conf.beat_schedule = { + "cleanup-stale-pending-tasks": { + "task": "worker.cleanup_stale_pending_tasks", + "schedule": 600.0, # 每 10 分钟(秒) + "options": {"expires": 300}, # 5 分钟过期,避免堆积 + }, +} diff --git a/apps/worker/worker_app/tasks/_startup.py b/apps/worker/worker_app/tasks/_startup.py index ff776c516..5c36e7fe3 100644 --- a/apps/worker/worker_app/tasks/_startup.py +++ b/apps/worker/worker_app/tasks/_startup.py @@ -10,6 +10,9 @@ logger = logging.getLogger(__name__) # 孤儿任务超时阈值:渲染任务超过此时间未更新则视为卡死 ORPHAN_TASK_TIMEOUT_MINUTES = 10 +# Pending 任务超时阈值:pending 任务在队列中等待超过此时间则自动清理 +PENDING_TASK_TIMEOUT_MINUTES = 30 + def cleanup_orphan_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover """清理数据库中超时未更新的 running GenerationTask(孤儿任务)。 @@ -83,6 +86,37 @@ def cleanup_stale_jobs(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> in return 0 +def cleanup_stale_pending_tasks(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover + """清理数据库中卡在 pending 状态超时的 GenerationTask。 + + 全局任务队列有 pending 数量上限,长期卡在 pending 的任务会占满队列, + 导致新用户无法创建任务。通过 created_at 超时判断并标记为 failed。 + + Args: + timeout_minutes: 超时时间(分钟),默认 PENDING_TASK_TIMEOUT_MINUTES + + Returns: + 清理的任务数量 + """ + from packages.adapters.sqlalchemy_impl.generation_task_repository import ( + SQLAlchemyGenerationTaskRepository, + ) + + try: + session = SessionLocal() + repo = SQLAlchemyGenerationTaskRepository(session) + count = repo.cleanup_stale_pending(timeout_minutes) + session.close() + if count > 0: + logger.warning("清理了 %d 个超时的 pending GenerationTask(超过 %d 分钟未处理)", count, timeout_minutes) + else: + logger.info("无超时 pending GenerationTask 需要清理") + return count + except Exception as e: + logger.error("清理超时 pending GenerationTask 失败: %s", e, exc_info=True) + return 0 + + def cleanup_all_stale_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> dict: # pragma: no cover """统一清理所有超时的孤儿任务。 @@ -93,15 +127,17 @@ def cleanup_all_stale_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) """ gen_count = cleanup_orphan_tasks(timeout_minutes) job_count = cleanup_stale_jobs(timeout_minutes) - total = gen_count + job_count + pending_count = cleanup_stale_pending_tasks(PENDING_TASK_TIMEOUT_MINUTES) + total = gen_count + job_count + pending_count if total > 0: logger.warning( - "孤儿任务清理完成: GenerationTask=%d, Job=%d, 总计=%d", + "任务清理完成: 孤儿 GenerationTask=%d, 孤儿 Job=%d, 超时 pending=%d, 总计=%d", gen_count, job_count, + pending_count, total, ) - return {"generation_tasks": gen_count, "jobs": job_count} + return {"generation_tasks": gen_count, "jobs": job_count, "pending": pending_count} @worker_ready.connect diff --git a/apps/worker/worker_app/tasks/cleanup.py b/apps/worker/worker_app/tasks/cleanup.py new file mode 100644 index 000000000..f3f6903d8 --- /dev/null +++ b/apps/worker/worker_app/tasks/cleanup.py @@ -0,0 +1,36 @@ +"""定期清理任务 — Celery Beat 调度。 + +包含: +- cleanup_stale_pending_tasks: 定期清理卡在 pending 超时的 generation_tasks +""" + +import logging + +from celery import shared_task + +from worker_app.tasks._startup import ( + PENDING_TASK_TIMEOUT_MINUTES, + cleanup_stale_pending_tasks, +) + +logger = logging.getLogger(__name__) + + +@shared_task(name="worker.cleanup_stale_pending_tasks") +def scheduled_cleanup_stale_pending(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINUTES) -> dict: + """Celery Beat 调度的定期任务:清理超时的 pending 任务。 + + 每 10 分钟执行一次(由 celery_app.py 的 beat_schedule 配置), + 查找所有 status='pending' 且 created_at < NOW() - timeout_minutes + 的 generation_tasks,批量更新为 failed。 + + Args: + timeout_minutes: 超时时间(分钟),默认 30 分钟 + + Returns: + {"cleaned": int} + """ + count = cleanup_stale_pending_tasks(timeout_minutes) + if count > 0: + logger.info("[Beat] 清理了 %d 个超时 pending 任务(超时阈值 %d 分钟)", count, timeout_minutes) + return {"cleaned": count} diff --git a/infra/docker/entrypoint-worker.sh b/infra/docker/entrypoint-worker.sh index d03eb4bd1..919bfc933 100755 --- a/infra/docker/entrypoint-worker.sh +++ b/infra/docker/entrypoint-worker.sh @@ -6,8 +6,11 @@ set -e CONCURRENCY="${WORKER_CONCURRENCY:-2}" +# 启用 -B 标志嵌入 celery beat(单 worker 实例,足够安全) +# beat 负责定期触发 pending 超时清理等定时任务 exec celery \ -A worker_app.celery_app \ worker \ --loglevel=info \ + "-B" \ "--concurrency=${CONCURRENCY}" diff --git a/packages/adapters/sqlalchemy_impl/generation_task_repository.py b/packages/adapters/sqlalchemy_impl/generation_task_repository.py index b0ed071af..fde48cd8a 100755 --- a/packages/adapters/sqlalchemy_impl/generation_task_repository.py +++ b/packages/adapters/sqlalchemy_impl/generation_task_repository.py @@ -310,3 +310,40 @@ class SQLAlchemyGenerationTaskRepository: model.completed_at = datetime.now(timezone.utc) self.session.commit() return len(models) + + def cleanup_stale_pending(self, timeout_minutes: int = 30) -> int: + """清理超时的 pending 任务(未被 Worker 拉取的任务)。 + + 全局任务队列有 pending 数量上限,长期卡在 pending 的任务会占满队列, + 导致新用户无法创建任务。将超时的 pending 任务标记为 failed。 + + Args: + timeout_minutes: 超时时间(分钟),默认 30 分钟 + + Returns: + 清理的任务数量 + """ + from datetime import timedelta + + cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes) + models = ( + self.session.query(GenerationTaskModel) + .filter( + GenerationTaskModel.status == GenerationTaskStatus.PENDING.value, + GenerationTaskModel.created_at < cutoff, + ) + .all() + ) + if not models: + return 0 + for model in models: + model.status = GenerationTaskStatus.FAILED.value + model.error_message = "pending timeout: auto cleanup" + model.error_info = { + "error_type": "PendingTimeout", + "message": f"任务在 pending 状态停留超过 {timeout_minutes} 分钟,自动清理", + "failed_at": datetime.now(timezone.utc).isoformat(), + } + model.completed_at = datetime.now(timezone.utc) + self.session.commit() + return len(models) diff --git a/tests/unit/test_generation_task_pending_cleanup.py b/tests/unit/test_generation_task_pending_cleanup.py new file mode 100644 index 000000000..ce257b937 --- /dev/null +++ b/tests/unit/test_generation_task_pending_cleanup.py @@ -0,0 +1,150 @@ +"""GenerationTaskRepository - cleanup_stale_pending 超时 pending 清理单元测试。""" + +import sys +from datetime import datetime, timedelta, timezone +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "apps" / "api")) + +from sqlalchemy import create_engine, text +from sqlalchemy.orm import sessionmaker + +from packages.adapters.sqlalchemy_impl.generation_task_repository import ( + SQLAlchemyGenerationTaskRepository, +) +from packages.adapters.sqlalchemy_impl.models import Base +from packages.domain import GenerationTask, GenerationTaskStatus + + +def _repository(): + engine = create_engine("sqlite:///:memory:") + Base.metadata.create_all(engine) + session = sessionmaker(bind=engine)() + return SQLAlchemyGenerationTaskRepository(session), session, engine + + +def _make_task(**kwargs) -> GenerationTask: + defaults = dict( + project_id="proj-1", + asset_library_id="lib-1", + created_by_user_id="user-1", + ) + defaults.update(kwargs) + return GenerationTask.create(**defaults) + + +# --------------------------------------------------------------------------- +# cleanup_stale_pending 基本测试 +# --------------------------------------------------------------------------- + +def test_cleanup_stale_pending_no_tasks_returns_zero(): + """没有任务时返回 0。""" + repo, _, _ = _repository() + count = repo.cleanup_stale_pending(timeout_minutes=30) + assert count == 0 + + +def test_cleanup_stale_pending_recent_pending_not_cleaned(): + """30 分钟内的 pending 任务不被清理。""" + repo, _, _ = _repository() + task = _make_task() + repo.create(task) + # 刚创建的 pending 任务不应被清理 + count = repo.cleanup_stale_pending(timeout_minutes=30) + assert count == 0 + assert repo.get(task.id).status == GenerationTaskStatus.PENDING + + +def test_cleanup_stale_pending_old_pending_marked_failed(): + """超过 30 分钟的 pending 任务被标记为 failed。""" + repo, _, engine = _repository() + task = _make_task() + repo.create(task) + + # 手动把 created_at 改到 1 小时前 + with engine.connect() as conn: + conn.execute( + text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"), + {"ts": datetime.now(timezone.utc) - timedelta(hours=1), "id": task.id}, + ) + conn.commit() + + count = repo.cleanup_stale_pending(timeout_minutes=30) + assert count == 1 + + saved = repo.get(task.id) + assert saved.status == GenerationTaskStatus.FAILED + assert saved.error_message == "pending timeout: auto cleanup" + assert saved.error_info.get("error_type") == "PendingTimeout" + assert "30" in saved.error_info["message"] + assert "failed_at" in saved.error_info + assert saved.completed_at is not None + + +def test_cleanup_stale_pending_running_not_touched(): + """running 任务不受影响,只清理 pending。""" + repo, _, engine = _repository() + task = _make_task() + repo.create(task) + task.mark_processing() + repo.update(task) + + # 回写 created_at 到 1 小时前 + with engine.connect() as conn: + conn.execute( + text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"), + {"ts": datetime.now(timezone.utc) - timedelta(hours=1), "id": task.id}, + ) + conn.commit() + + count = repo.cleanup_stale_pending(timeout_minutes=30) + assert count == 0 + assert repo.get(task.id).status == GenerationTaskStatus.RUNNING + + +def test_cleanup_stale_pending_custom_timeout(): + """自定义超时时间生效。""" + repo, _, engine = _repository() + task = _make_task() + repo.create(task) + + # 回写 created_at 到 20 分钟前 + with engine.connect() as conn: + conn.execute( + text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"), + {"ts": datetime.now(timezone.utc) - timedelta(minutes=20), "id": task.id}, + ) + conn.commit() + + # 30 分钟超时:不清理 + count_30 = repo.cleanup_stale_pending(timeout_minutes=30) + assert count_30 == 0 + # 15 分钟超时:清理 + count_15 = repo.cleanup_stale_pending(timeout_minutes=15) + assert count_15 == 1 + assert repo.get(task.id).status == GenerationTaskStatus.FAILED + + +def test_cleanup_stale_pending_multiple(): + """批量清理多个超时的 pending 任务。""" + repo, _, engine = _repository() + + tasks = [] + for i in range(5): + t = _make_task(project_id=f"proj-{i}") + repo.create(t) + tasks.append(t) + + # 全部回写 created_at 到 2 小时前 + with engine.connect() as conn: + for t in tasks: + conn.execute( + text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"), + {"ts": datetime.now(timezone.utc) - timedelta(hours=2), "id": t.id}, + ) + conn.commit() + + count = repo.cleanup_stale_pending(timeout_minutes=30) + assert count == 5 + for t in tasks: + assert repo.get(t.id).status == GenerationTaskStatus.FAILED -- 2.54.0 From d3a9fff0c211c6b3a41bbb686aed8504c93eaae3 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sun, 23 Aug 2026 04:24:58 +0000 Subject: [PATCH 2/4] style: auto-format with black + isort + prettier [skip ci-format-check] --- apps/worker/worker_app/tasks/cleanup.py | 1 - tests/unit/test_generation_task_pending_cleanup.py | 1 + 2 files changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/worker/worker_app/tasks/cleanup.py b/apps/worker/worker_app/tasks/cleanup.py index f3f6903d8..28d5e7e22 100644 --- a/apps/worker/worker_app/tasks/cleanup.py +++ b/apps/worker/worker_app/tasks/cleanup.py @@ -7,7 +7,6 @@ import logging from celery import shared_task - from worker_app.tasks._startup import ( PENDING_TASK_TIMEOUT_MINUTES, cleanup_stale_pending_tasks, diff --git a/tests/unit/test_generation_task_pending_cleanup.py b/tests/unit/test_generation_task_pending_cleanup.py index ce257b937..252f1be19 100644 --- a/tests/unit/test_generation_task_pending_cleanup.py +++ b/tests/unit/test_generation_task_pending_cleanup.py @@ -37,6 +37,7 @@ def _make_task(**kwargs) -> GenerationTask: # cleanup_stale_pending 基本测试 # --------------------------------------------------------------------------- + def test_cleanup_stale_pending_no_tasks_returns_zero(): """没有任务时返回 0。""" repo, _, _ = _repository() -- 2.54.0 From 1712c1693af0ec556f8d36108957af89b7109e87 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Sun, 23 Aug 2026 12:39:49 +0800 Subject: [PATCH 3/4] fix: address AI review - session leak & bulk update for pending cleanup - _startup.py: move session.close() to finally block to prevent DB connection leak - generation_task_repository.py: replace .all()+loop with bulk .update() to avoid OOM when large number of pending tasks accumulate --- apps/worker/worker_app/tasks/_startup.py | 5 ++-- .../generation_task_repository.py | 30 ++++++++++--------- 2 files changed, 19 insertions(+), 16 deletions(-) diff --git a/apps/worker/worker_app/tasks/_startup.py b/apps/worker/worker_app/tasks/_startup.py index 5c36e7fe3..c841f4406 100644 --- a/apps/worker/worker_app/tasks/_startup.py +++ b/apps/worker/worker_app/tasks/_startup.py @@ -102,11 +102,10 @@ def cleanup_stale_pending_tasks(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINU SQLAlchemyGenerationTaskRepository, ) + session = SessionLocal() try: - session = SessionLocal() repo = SQLAlchemyGenerationTaskRepository(session) count = repo.cleanup_stale_pending(timeout_minutes) - session.close() if count > 0: logger.warning("清理了 %d 个超时的 pending GenerationTask(超过 %d 分钟未处理)", count, timeout_minutes) else: @@ -115,6 +114,8 @@ def cleanup_stale_pending_tasks(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINU except Exception as e: logger.error("清理超时 pending GenerationTask 失败: %s", e, exc_info=True) return 0 + finally: + session.close() def cleanup_all_stale_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> dict: # pragma: no cover diff --git a/packages/adapters/sqlalchemy_impl/generation_task_repository.py b/packages/adapters/sqlalchemy_impl/generation_task_repository.py index fde48cd8a..b9d318672 100755 --- a/packages/adapters/sqlalchemy_impl/generation_task_repository.py +++ b/packages/adapters/sqlalchemy_impl/generation_task_repository.py @@ -326,24 +326,26 @@ class SQLAlchemyGenerationTaskRepository: from datetime import timedelta cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes) - models = ( + error_info = { + "error_type": "PendingTimeout", + "message": f"任务在 pending 状态停留超过 {timeout_minutes} 分钟,自动清理", + "failed_at": datetime.now(timezone.utc).isoformat(), + } + count = ( self.session.query(GenerationTaskModel) .filter( GenerationTaskModel.status == GenerationTaskStatus.PENDING.value, GenerationTaskModel.created_at < cutoff, ) - .all() + .update( + { + GenerationTaskModel.status: GenerationTaskStatus.FAILED.value, + GenerationTaskModel.error_message: "pending timeout: auto cleanup", + GenerationTaskModel.error_info: error_info, + GenerationTaskModel.completed_at: datetime.now(timezone.utc), + }, + synchronize_session="fetch", + ) ) - if not models: - return 0 - for model in models: - model.status = GenerationTaskStatus.FAILED.value - model.error_message = "pending timeout: auto cleanup" - model.error_info = { - "error_type": "PendingTimeout", - "message": f"任务在 pending 状态停留超过 {timeout_minutes} 分钟,自动清理", - "failed_at": datetime.now(timezone.utc).isoformat(), - } - model.completed_at = datetime.now(timezone.utc) self.session.commit() - return len(models) + return count -- 2.54.0 From 498b490bac2edc8e1603fdfb29f50fd4f9bc83c9 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Sun, 23 Aug 2026 12:44:50 +0800 Subject: [PATCH 4/4] fix: document single-instance constraint for embedded Beat & optimize synchronize_session - entrypoint-worker.sh: prominent deployment constraint warning (replicas=1 only) - generation_task_repository.py: synchronize_session=False for cleanup (no post-update session sync needed) --- infra/docker/entrypoint-worker.sh | 6 ++++-- .../adapters/sqlalchemy_impl/generation_task_repository.py | 2 +- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/infra/docker/entrypoint-worker.sh b/infra/docker/entrypoint-worker.sh index 919bfc933..f2e208d96 100755 --- a/infra/docker/entrypoint-worker.sh +++ b/infra/docker/entrypoint-worker.sh @@ -6,8 +6,10 @@ set -e CONCURRENCY="${WORKER_CONCURRENCY:-2}" -# 启用 -B 标志嵌入 celery beat(单 worker 实例,足够安全) -# beat 负责定期触发 pending 超时清理等定时任务 +# ⚠️ 部署约束:此 Worker 必须且只能运行单实例(replicas=1) +# -B 标志嵌入 celery beat,beat 负责定期触发 pending 超时清理等定时任务 +# 多实例部署会导致每个 Worker 独立运行 Beat,造成定时任务重复执行 +# 若需横向扩展 Worker,必须将 Beat 拆分为独立服务(celery beat -A worker_app.celery_app) exec celery \ -A worker_app.celery_app \ worker \ diff --git a/packages/adapters/sqlalchemy_impl/generation_task_repository.py b/packages/adapters/sqlalchemy_impl/generation_task_repository.py index b9d318672..75f394edc 100755 --- a/packages/adapters/sqlalchemy_impl/generation_task_repository.py +++ b/packages/adapters/sqlalchemy_impl/generation_task_repository.py @@ -344,7 +344,7 @@ class SQLAlchemyGenerationTaskRepository: GenerationTaskModel.error_info: error_info, GenerationTaskModel.completed_at: datetime.now(timezone.utc), }, - synchronize_session="fetch", + synchronize_session=False, ) ) self.session.commit() -- 2.54.0