fix(queue): #2073 Worker 队列分流——任务路由补全 + beat 独立 + transcode 并发独立 #2079
Reference in New Issue
Block a user
Delete Branch "fix/worker-queue-split"
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?
背景
Staging 出现「用户没发多任务却排队」现象。排查发现:
packages/shared/celery_queues.py的task_routes只显式路由了 5 个任务,大量实时任务(TTS、音色克隆、人声/背景提取、AI 数字人、GPU MuseTalk lipsync 全链路)没配路由,默认落到celery队列被 transcode worker 消费;generation worker 根本收不到这些任务 → 实时任务被后台队列阻塞。entrypoint-worker.sh里 celery beat 嵌在 generation worker 里(-B),beat 进程本身不占执行槽但日志/生命周期耦合,且启动并发逻辑中 TRANSCODE_CONCURRENCY 依赖WORKER_CONCURRENCY - GENERATION_CONCURRENCY差值(WORKER_CONCURRENCY=2, GENERATION_CONCURRENCY=2 时 trans=0→硬编码兜底 1,运维无法独立调)。TRANSCODE_CONCURRENCY/ beat 开关,资源限制与现状不符。97ad0ae2时期calculate_asset_quality因'str'.valuebug 连续失败,8abdeb95 已修复但历史素材 quality_score 仍为 NULL,需要一次性补跑。变更
(a) 任务路由补全(核心修复)
packages/shared/celery_queues.py所有实时链路任务显式路由到
generation队列:worker.generate_videoworker.process_tts_synthesis/worker.process_tts_segment_synthesisworker.process_voice_cloneworker.extract_voice/worker.extract_backgroundai_avatar_render.executelipsync_gpu_process_async/lipsync_tts.synthesize_and_submit/lipsync_tts.poll_mediakit_status/lipsync_tts.persist_output_videoworker.ingest_assetworker.classify_assetworker.calculate_asset_qualityworker.generate_atom_clips/worker.tag_atom_clip/worker.backfill_atom_clip_tagsworker.process_duplication_check/worker.check_duplicateworker.batch_download_videos/worker.batch_generate_thumbnailsworker.cleanup_stale_*(4 个)(b) Transcode 并发独立可配
infra/docker/entrypoint-worker.shTRANSCODE_CONCURRENCY默认 2(不再用WORKER_CONCURRENCY - GENERATION_CONCURRENCY差值)GENERATION_CONCURRENCY默认 2(保持)*_CONCURRENCY都未显式设置时,才用WORKER_CONCURRENCY按比例对半分配(兼容旧配置)WORKER_MAX_TASKS_PER_CHILD默认 100(与 compose.yml 一致)(c) Beat 独立进程
infra/docker/entrypoint-worker.sh-Q generation)+ transcode worker(-Q transcode,celery)-B,beat 不再与 generation 生命周期耦合BEAT_ENABLED=1/0开关,未来独立 beat 容器部署时 worker 容器设为 0 即可(d) 资源限制与 compose
infra/docker/compose.ymlGENERATION_CONCURRENCY/TRANSCODE_CONCURRENCY/BEAT_ENABLED三个环境变量.env.example同步补全三个新变量的说明(e) 历史 quality 补打分脚本
apps/worker/scripts/backfill_asset_quality.py一次性脚本(不常驻注册到 celery imports),在 worker 容器内执行:
不做的事
task_enqueue.py的WORKER_CONCURRENCY = 4限流阈值常量(这是业务限流估算用,不影响实际 celery 并发;留待后续与前端排队估算一并调)。部署步骤(staging)
.env里按需显式设置(不设则用默认 2/2/1):docker exec xiaoxia-worker-staging ps -ef | grep celery→ 应看到 3 个 celery 进程(beat + generation worker + transcode worker)docker exec xiaoxia-worker-staging celery -A worker_app.celery_app inspect active -d generation@<host>和-d transcode@<host>分别确认两个 worker 在跑.env里运维热调的GENERATION_CONCURRENCY=2保留,新增显式TRANSCODE_CONCURRENCY=2。97ad0ae2时期 calculate_asset_quality 因 str.value bug 失败,8abdeb95 已修复; 本脚本扫描 quality_score IS NULL 的视频素材投递到 transcode 队列,支持 dry-run/时间范围/分批限流🚀 预览环境已部署
🗑️ 预览环境已清理
PR #2079 已关闭或合并,对应的预览环境已被清理。