From 7ed0ffd5a8acabef825585d87c88c9f2dcaf3ecc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=94=A8=E6=88=B7CI=20Test?= Date: Sat, 11 Jul 2026 01:15:21 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E7=94=9F=E6=88=90=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E5=85=A5=E9=98=9F=E5=A4=B1=E8=B4=A5=E6=97=B6=E6=A0=87=E8=AE=B0?= =?UTF-8?q?=E4=B8=BAfailed=EF=BC=8C=E9=81=BF=E5=85=8Dpending=E5=83=B5?= =?UTF-8?q?=E5=B0=B8=E4=BB=BB=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因:任务创建(DB commit)和入队(Celery send_task)是两个独立操作, send_task失败时任务卡在pending状态永远不会执行。 修复: - 新增_safe_enqueue_generation_task安全入队函数 - send_task失败时自动标记任务为failed并记录错误 - 覆盖4处入口:批量创建、generation重试、task_center两级重试 --- apps/api/app/api/routes/generation_tasks.py | 55 +++++++++++++++++---- apps/api/app/api/routes/task_center.py | 35 ++++++++++++- 2 files changed, 78 insertions(+), 12 deletions(-) mode change 100644 => 100755 apps/api/app/api/routes/generation_tasks.py mode change 100644 => 100755 apps/api/app/api/routes/task_center.py diff --git a/apps/api/app/api/routes/generation_tasks.py b/apps/api/app/api/routes/generation_tasks.py old mode 100644 new mode 100755 index 36c11f709..ebedda705 --- a/apps/api/app/api/routes/generation_tasks.py +++ b/apps/api/app/api/routes/generation_tasks.py @@ -37,6 +37,43 @@ logger = logging.getLogger(__name__) router = APIRouter() +def _safe_enqueue_generation_task( + task: Any, + generation_task_repository: Any, +) -> bool: + """安全入队:send_task 失败时自动把任务标记为 failed,避免留下 pending 僵尸任务。 + + Returns: + True 表示入队成功,False 表示入队失败(已标记为 failed) + """ + try: + celery_app.send_task("worker.generate_video", args=[task.id]) + logger.info( + "[生成任务] 入队成功: task_id=%s, status=%s", + task.id, + task.status, + ) + return True + except Exception as e: + logger.error( + "[生成任务] 入队失败,标记为失败: task_id=%s error=%s", + task.id, + e, + exc_info=True, + ) + try: + task.mark_failed(f"任务入队失败: {e}") + generation_task_repository.update(task) + except Exception as update_err: + logger.error( + "[生成任务] 入队失败后更新状态也失败: task_id=%s error=%s", + task.id, + update_err, + exc_info=True, + ) + return False + + def _check_project_access(project_id: str, user_id: str, project_repository) -> None: """检查用户是否有项目访问权限""" project = project_repository.find_by_id(project_id) @@ -228,6 +265,7 @@ def create_generation_task( use_case = CreateGenerationTaskUseCase(generation_task_repository) count = request.count created_tasks = [] + failed_tasks = [] # 同批次任务共享 batch_id,用于视频查重时批次内比对 batch_id = uuid.uuid4().hex if count > 1 else "" @@ -249,19 +287,15 @@ def create_generation_task( batch_id=batch_id, ) ) - celery_app.send_task("worker.generate_video", args=[task.id]) - created_tasks.append(task) - logger.info( - "[生成任务] 入队成功: task_id=%s, status=%s, batch_id=%s", - task.id, - task.status, - batch_id, - ) + if _safe_enqueue_generation_task(task, generation_task_repository): + created_tasks.append(task) + else: + failed_tasks.append(task) except Exception as e: logger.error("[生成任务] 创建失败: %s", e, exc_info=True) raise HTTPException(status_code=500, detail="创建生成任务失败,请稍后重试或查看任务日志") - items = [_to_generation_task_response(t) for t in created_tasks] + items = [_to_generation_task_response(t) for t in created_tasks + failed_tasks] return BatchGenerationTaskResponse(items=items, total=len(items)) @@ -347,5 +381,6 @@ def retry_generation_task( asset_select_mode=getattr(task, "asset_select_mode", ""), ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[生成任务] 重试入队失败: task_id=%s", retried.id) return _to_generation_task_response(retried) diff --git a/apps/api/app/api/routes/task_center.py b/apps/api/app/api/routes/task_center.py old mode 100644 new mode 100755 index 62c97c960..13e51f480 --- a/apps/api/app/api/routes/task_center.py +++ b/apps/api/app/api/routes/task_center.py @@ -25,6 +25,35 @@ from packages.application import ( router = APIRouter() +def _safe_enqueue_generation_task( + task: Any, + generation_task_repository: Any, +) -> bool: + """安全入队:send_task 失败时自动把任务标记为 failed,避免留下 pending 僵尸任务。""" + try: + celery_app.send_task("worker.generate_video", args=[task.id]) + logger.info("[任务中心] 生成任务入队成功: task_id=%s", task.id) + return True + except Exception as e: + logger.error( + "[任务中心] 生成任务入队失败,标记为失败: task_id=%s error=%s", + task.id, + e, + exc_info=True, + ) + try: + task.mark_failed(f"任务入队失败: {e}") + generation_task_repository.update(task) + except Exception as update_err: + logger.error( + "[任务中心] 入队失败后更新状态也失败: task_id=%s error=%s", + task.id, + update_err, + exc_info=True, + ) + return False + + def _humanize_task_error(error_message: str) -> str: raw = (error_message or "").strip() if not raw: @@ -153,7 +182,8 @@ def retry_task_by_id( created_by_user_id=authenticated_user.user.id, ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[任务中心] 用户级重试入队失败: task_id=%s", retried.id) return UserTaskResponse( id=f"generation:{retried.id}", task_type="generation", @@ -235,7 +265,8 @@ def retry_project_task( created_by_user_id=authenticated_user.user.id, ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[任务中心] 项目级重试用队失败: task_id=%s", retried.id) return _generation_task_to_project_response(retried) if task_type == "ingest": job = ingest_job_repository.get(source_id) -- 2.54.0