feat(worker): celery 队列隔离 + 孤儿任务消息作废 (#1714) #1722
Reference in New Issue
Block a user
Delete Branch "feature/celery-queue-isolation-1714"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
背景(Issue #1714 补充任务)
ingest_asset与视频生成共用默认队列、worker 单进程消费,实测 21 个转码积压把用户生成任务堵 40+ 分钟。failed → running非法状态转换,worker 打 ERROR 后仍跑完产出半成品。改动
队列隔离(三队列拓扑)
packages/shared/celery_queues.py:generation/transcode/celery三队列 +task_routesworker.generate_video→ generation(用户视频生成,高优先级)ingest_asset/classify_asset/process_duplication_check/check_duplicate→ transcode(素材入库后处理链路)entrypoint-worker.sh改双进程:-Q generation -B(beat 内嵌,单实例约束),worker_prefetch_multiplier=1、task_acks_late=True,并发GENERATION_CONCURRENCY默认 2-Q transcode,celery,并发=总并发-2(最小 1)wait -n任一退出则 TERM 另一个;compose / deploy-staging / deploy-production / start-worker.ps1 同步消息作废(双保险)
packages/shared/celery_orphan_guard.py:revoke_and_purge():control.revoke 广播(在线 worker)+ 直接扫 Redis list 物理移除消息体(worker 下线期间 revoke 收不到);解析 celery JSON 信封,按业务 id + celery headers.id 双匹配,未命中 rpush 回写保序,解析失败保守保留ensure_task_claimable():终态(failed/cancelled/completed)抛 StaleTaskDiscardedgeneration_tasks.celery_task_id/ingest_jobs.celery_task_id(新列,067 迁移,server_default 空串,写入失败仅 warning 不阻断);覆盖生成/上传/分片/重试全部入队点{"status":"discarded"}不进业务逻辑;mark_processing返回 False(非法转换如 failed→running)安全中止,不再打 ERROR 后继续渲染测试
staging 验证方式
部署后:
docker exec xiaoxia-worker-staging pgrep -af celery应见 generation@/transcode@ 两进程;celery inspect active_queues 确认消费队列分离;向 transcode 灌入 20+ 任务后提交生成任务,1 分钟内开始渲染。配套
🚀 预览环境已部署
f5bdf17da1todf99305dd6【阻塞级判定】
📊 审查概览
🔴 阻塞级问题(必须修复)
[apps/api/app/api/routes/chunked_upload.py:384] 任务ID持久化失败导致孤儿任务无法清理
celery_task_id持久化到数据库,但使用了裸except Exception: pass。如果数据库更新失败,celery_task_id将无法记录。当后续触发超时清理逻辑时,_revoke_and_purge_stale_messages将无法获取到有效的celery_task_id,导致无法撤销队列中的消息。这会导致DB中任务状态已标记为失败,但队列中仍存在旧消息,Worker重启后可能重新消费并执行,造成非法状态转换和资源浪费。celery_app.control.revoke撤销已发送的任务。[apps/api/app/api/routes/ingest_jobs.py:46] 任务ID持久化失败导致孤儿任务无法清理
ingest_jobs.py中使用了相同的裸except Exception: pass模式处理celery_task_id的持久化。这会导致素材转码任务在DB更新失败时丢失与Celery消息的关联,后续清理流程无法正确撤销消息,可能引发重复处理。[apps/api/app/api/routes/task_center.py:378] 任务ID持久化失败导致孤儿任务无法清理
task_center.py中的重试逻辑也存在相同问题。静默忽略持久化错误会导致重试生成的任务与Celery消息断链,破坏了#1714 issue修复的核心目标(通过ID关联进行消息清理)。💡 改进建议(不阻塞合并)
[apps/worker/worker_app/tasks/_startup.py:45] 降级逻辑可能掩盖真实错误
cleanup_stale_running_with_session_ids中,当repo不支持_with_ids方法时,降级返回[("", "")]。虽然这兼容了旧代码,但如果cleanup_stale_running本身抛出异常(如数据库连接问题),这里会直接崩溃,且没有日志提示“使用了降级模式”。建议在降级分支增加日志记录,方便排查问题。[packages/adapters/sqlalchemy_impl/generation_task_repository.py:329] 批量更新改为循环更新影响性能
cleanup_stale_pending_with_ids方法中,原代码使用了.update(..., synchronize_session=False)进行批量更新,效率较高。修改后的代码先.all()查出所有对象,再循环修改属性并提交。如果待清理的pending任务数量较多(例如积压场景),这会一次性加载大量数据到内存,且产生大量的UPDATE语句。建议在数据量可控的情况下保留,或者考虑分批处理以平衡内存和性能。✅ 良好实践
generation.py和ingest.py)增加了执行前的状态守卫检查,有效防止了僵尸消息的非法执行,这是解决并发/分布式环境下状态一致性的关键措施。celery_task_id设置了server_default="",保证了非空约束下的平滑升级。entrypoint-worker.sh脚本中通过trap和wait -n实现了双Worker进程的协同退出管理,避免了孤儿进程。✅ 格式检查通过 | ❌ 逻辑审查需修改 | ⚠️ 建议关注性能
🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 11:14:26 | 模型:
🗑️ 预览环境已清理
PR #1722 已关闭或合并,对应的预览环境已被清理。