Compare commits

...

1 Commits

Author SHA1 Message Date
用户CI Test 7ed0ffd5a8 fix: 生成任务入队失败时标记为failed,避免pending僵尸任务
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 8s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m0s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
根因:任务创建(DB commit)和入队(Celery send_task)是两个独立操作,
send_task失败时任务卡在pending状态永远不会执行。

修复:
- 新增_safe_enqueue_generation_task安全入队函数
- send_task失败时自动标记任务为failed并记录错误
- 覆盖4处入口:批量创建、generation重试、task_center两级重试
2026-07-11 01:15:21 +08:00
2 changed files with 78 additions and 12 deletions
+45 -10
View File
@@ -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)
+33 -2
View File
@@ -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)