feat(worker): celery 队列隔离 + 孤儿任务消息作废 (#1714)
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m27s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m39s
AI Code Review / AI Code Review (pull_request) Successful in 4m18s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 11m21s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 14s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m38s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 1m49s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 3m1s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m26s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 5m3s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 6m19s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 260h20m29s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 260h20m33s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 260h20m33s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 260h26m40s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 260h26m47s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 260h26m40s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 260h26m47s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 260h26m40s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 260h26m53s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 260h26m54s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 260h26m53s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 260h27m0s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 260h55m7s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 261h1m16s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 261h1m21s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 261h1m27s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 261h1m28s
CI/CD Pipeline / Deploy Production (pull_request) Failing after 260h55m3s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 261h1m27s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m27s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m39s
AI Code Review / AI Code Review (pull_request) Successful in 4m18s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 11m21s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 14s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m38s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 1m49s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 3m1s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m26s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 5m3s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 6m19s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 260h20m29s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 260h20m33s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 260h20m33s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 260h26m40s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 260h26m47s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 260h26m40s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 260h26m47s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 260h26m40s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 260h26m53s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 260h26m54s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 260h26m53s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 260h27m0s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 260h55m7s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 261h1m16s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 261h1m21s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 261h1m27s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 261h1m28s
CI/CD Pipeline / Deploy Production (pull_request) Failing after 260h55m3s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 261h1m27s
问题:素材转码与视频生成共用 celery 默认队列、worker 单进程消费, 20+ 转码积压会把用户生成任务堵 40 分钟以上;孤儿清理把任务标 failed 后 Redis 队列消息未作废,消息被重投导致 failed→running 非法转换, worker 打印 ERROR 后继续产出半成品。 队列隔离: - 新增 packages/shared/celery_queues.py:generation/transcode/celery 三队列与 task_routes(generate_video→generation;ingest_asset/ classify_asset/duplication→transcode),apply_queue_settings() - worker 入口改双进程:generation worker 独占队列并内嵌 beat (prefetch=1, GENERATION_CONCURRENCY 默认 2),transcode worker 消费 transcode,celery(并发=总-2,最小 1),任一退出则整体终止 - compose/部署脚本/ps1 同步新增 GENERATION_CONCURRENCY 与健康检查 消息作废: - 新增 packages/shared/celery_orphan_guard.py:终态守卫 ensure_task_claimable、Redis 队列消息物理清理(JSON 信封解析, 按业务 id + celery headers.id 双匹配,未命中 rpush 保序)、 revoke_and_purge(control.revoke + 物理清队列双保险) - 入队点(生成/上传/分片/重试)send_task 后持久化 celery_task_id 到 generation_tasks/ingest_jobs(新列,067 迁移,失败仅 warning) - generate_video/ingest_asset 执行前校验 DB 状态:终态直接 discarded 不进业务逻辑;mark_processing 返回 False(非法转换)安全中止 - 孤儿/超时清理标 failed 时同时 revoke + 清队列消息 - pending 超时阈值 15→45 分钟,与 running 孤儿(20min)区分 测试:新增 22 个单测(路由表/真实 Redis 消息清理/终态守卫/ 非法转换中止/标 failed 后消息不重投/入队持久化),全量 14301 passed;067 迁移隔离 DDL 验证 upgrade/downgrade 通过。
This commit is contained in:
@@ -38,6 +38,7 @@ def _to_domain(model: GenerationTaskModel) -> GenerationTask:
|
||||
bgm_config=dict(getattr(model, "bgm_config", {}) or {}),
|
||||
is_preview=bool(getattr(model, "is_preview", False)),
|
||||
source_task_id=getattr(model, "source_task_id", "") or "",
|
||||
celery_task_id=getattr(model, "celery_task_id", "") or "",
|
||||
output_width=getattr(model, "output_width", 1280) or 1280,
|
||||
output_height=getattr(model, "output_height", 720) or 720,
|
||||
cover_url=getattr(model, "cover_url", "") or "",
|
||||
@@ -82,6 +83,7 @@ class SQLAlchemyGenerationTaskRepository:
|
||||
bgm_config=task.bgm_config or {},
|
||||
is_preview=task.is_preview or False,
|
||||
source_task_id=task.source_task_id or "",
|
||||
celery_task_id=getattr(task, "celery_task_id", "") or "",
|
||||
output_width=task.output_width,
|
||||
output_height=task.output_height,
|
||||
cover_url=task.cover_url or "",
|
||||
@@ -315,6 +317,7 @@ class SQLAlchemyGenerationTaskRepository:
|
||||
if hasattr(model, "is_preview"):
|
||||
model.is_preview = task.is_preview or False
|
||||
model.source_task_id = task.source_task_id or ""
|
||||
model.celery_task_id = getattr(task, "celery_task_id", "") or model.celery_task_id or ""
|
||||
model.output_width = task.output_width
|
||||
model.output_height = task.output_height
|
||||
model.cover_url = task.cover_url or ""
|
||||
@@ -326,12 +329,14 @@ class SQLAlchemyGenerationTaskRepository:
|
||||
def cleanup_stale_running(self, timeout_minutes: int = 10) -> int:
|
||||
"""清理超时未更新的 running 任务(孤儿任务)。
|
||||
|
||||
将 status=running 且 updated_at 超过 timeout_minutes 分钟未更新的任务
|
||||
标记为 failed,error_message 标记为任务执行中断。
|
||||
|
||||
Returns:
|
||||
清理的任务数量
|
||||
清理的任务数量(仅计数,保持旧签名兼容)
|
||||
"""
|
||||
items = self.cleanup_stale_running_with_ids(timeout_minutes)
|
||||
return len(items)
|
||||
|
||||
def cleanup_stale_running_with_ids(self, timeout_minutes: int = 10) -> list[tuple[str, str]]:
|
||||
"""同 cleanup_stale_running,但返回 [(task_id, celery_task_id), ...] 供撤销队列消息。"""
|
||||
from datetime import timedelta
|
||||
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes)
|
||||
@@ -344,8 +349,10 @@ class SQLAlchemyGenerationTaskRepository:
|
||||
.all()
|
||||
)
|
||||
if not models:
|
||||
return 0
|
||||
return []
|
||||
result: list[tuple[str, str]] = []
|
||||
for model in models:
|
||||
result.append((model.id, getattr(model, "celery_task_id", "") or ""))
|
||||
model.status = GenerationTaskStatus.FAILED.value
|
||||
model.error_message = "任务执行中断(worker重启/超时)"
|
||||
model.error_info = {
|
||||
@@ -355,43 +362,43 @@ class SQLAlchemyGenerationTaskRepository:
|
||||
}
|
||||
model.completed_at = datetime.now(timezone.utc)
|
||||
self.session.commit()
|
||||
return len(models)
|
||||
return result
|
||||
|
||||
def cleanup_stale_pending(self, timeout_minutes: int = 30) -> int:
|
||||
"""清理超时的 pending 任务(未被 Worker 拉取的任务)。
|
||||
|
||||
全局任务队列有 pending 数量上限,长期卡在 pending 的任务会占满队列,
|
||||
导致新用户无法创建任务。将超时的 pending 任务标记为 failed。
|
||||
|
||||
Args:
|
||||
timeout_minutes: 超时时间(分钟),默认 30 分钟
|
||||
|
||||
Returns:
|
||||
清理的任务数量
|
||||
清理的任务数量(仅计数,保持旧签名兼容)
|
||||
"""
|
||||
items = self.cleanup_stale_pending_with_ids(timeout_minutes)
|
||||
return len(items)
|
||||
|
||||
def cleanup_stale_pending_with_ids(self, timeout_minutes: int = 30) -> list[tuple[str, str]]:
|
||||
"""同 cleanup_stale_pending,但返回 [(task_id, celery_task_id), ...] 供撤销队列消息。"""
|
||||
from datetime import timedelta
|
||||
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes)
|
||||
error_info = {
|
||||
"error_type": "PendingTimeout",
|
||||
"message": f"任务在 pending 状态停留超过 {timeout_minutes} 分钟,自动清理",
|
||||
"failed_at": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
count = (
|
||||
models = (
|
||||
self.session.query(GenerationTaskModel)
|
||||
.filter(
|
||||
GenerationTaskModel.status == GenerationTaskStatus.PENDING.value,
|
||||
GenerationTaskModel.created_at < cutoff,
|
||||
)
|
||||
.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=False,
|
||||
)
|
||||
.all()
|
||||
)
|
||||
if not models:
|
||||
return []
|
||||
error_info = {
|
||||
"error_type": "PendingTimeout",
|
||||
"message": f"任务在 pending 状态停留超过 {timeout_minutes} 分钟,自动清理",
|
||||
"failed_at": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
result: list[tuple[str, str]] = []
|
||||
for model in models:
|
||||
result.append((model.id, getattr(model, "celery_task_id", "") or ""))
|
||||
model.status = GenerationTaskStatus.FAILED.value
|
||||
model.error_message = "pending timeout: auto cleanup"
|
||||
model.error_info = error_info
|
||||
model.completed_at = datetime.now(timezone.utc)
|
||||
self.session.commit()
|
||||
return count
|
||||
return result
|
||||
|
||||
@@ -19,6 +19,7 @@ class SQLAlchemyIngestJobRepository:
|
||||
result_asset_id=job.result_asset_id,
|
||||
file_hash=job.file_hash,
|
||||
asset_id=job.asset_id or "",
|
||||
celery_task_id=getattr(job, "celery_task_id", "") or "",
|
||||
created_at=job.created_at,
|
||||
updated_at=job.updated_at,
|
||||
)
|
||||
@@ -40,6 +41,7 @@ class SQLAlchemyIngestJobRepository:
|
||||
result_asset_id=model.result_asset_id,
|
||||
file_hash=model.file_hash or "",
|
||||
asset_id=getattr(model, "asset_id", "") or "",
|
||||
celery_task_id=getattr(model, "celery_task_id", "") or "",
|
||||
created_at=model.created_at,
|
||||
updated_at=model.updated_at,
|
||||
)
|
||||
@@ -59,6 +61,9 @@ class SQLAlchemyIngestJobRepository:
|
||||
model.storage_key = job.storage_key
|
||||
if job.asset_id:
|
||||
model.asset_id = job.asset_id
|
||||
celery_tid = getattr(job, "celery_task_id", "")
|
||||
if celery_tid:
|
||||
model.celery_task_id = celery_tid
|
||||
model.updated_at = job.updated_at
|
||||
self.session.commit()
|
||||
return job
|
||||
|
||||
@@ -246,6 +246,7 @@ class IngestJobModel(Base):
|
||||
result_asset_id = Column(String(36), nullable=False, default="")
|
||||
file_hash = Column(String(64), nullable=True, index=True)
|
||||
asset_id = Column(String(36), nullable=False, default="", index=True)
|
||||
celery_task_id = Column(String(64), nullable=False, default="", server_default="")
|
||||
created_at = Column(DateTime, nullable=False, default=lambda: datetime.now(timezone.utc))
|
||||
updated_at = Column(DateTime, nullable=False, default=lambda: datetime.now(timezone.utc))
|
||||
|
||||
@@ -298,6 +299,7 @@ class GenerationTaskModel(Base):
|
||||
resolution = Column(String(20), nullable=False, default="")
|
||||
is_preview = Column(Boolean, nullable=False, default=False, index=True)
|
||||
source_task_id = Column(String(32), nullable=False, default="", index=True)
|
||||
celery_task_id = Column(String(64), nullable=False, default="", server_default="")
|
||||
output_width = Column(Integer, nullable=False, default=1280)
|
||||
output_height = Column(Integer, nullable=False, default=720)
|
||||
cover_url = Column(String(1000), nullable=False, default="")
|
||||
|
||||
Reference in New Issue
Block a user