fix(queue): #2073 Worker 队列分流——任务路由补全 + beat 独立 + transcode 并发独立 #2079

Merged
auto-approve-bot merged 7 commits from fix/worker-queue-split into develop 2026-09-28 01:09:31 +08:00
Owner

背景

Staging 出现「用户没发多任务却排队」现象。排查发现:

  1. packages/shared/celery_queues.py 的 task_routes 只显式路由了 5 个任务,大量实时任务(TTS、音色克隆、人声/背景提取、AI 数字人、GPU MuseTalk lipsync 全链路)没配路由,默认落到 celery 队列被 transcode worker 消费;generation worker 根本收不到这些任务 → 实时任务被后台队列阻塞。
  2. entrypoint-worker.sh 里 celery beat 嵌在 generation worker 里(-B),beat 进程本身不占执行槽但日志/生命周期耦合,且启动并发逻辑中 TRANSCODE_CONCURRENCY 依赖 WORKER_CONCURRENCY - GENERATION_CONCURRENCY 差值(WORKER_CONCURRENCY=2, GENERATION_CONCURRENCY=2 时 trans=0→硬编码兜底 1,运维无法独立调)。
  3. compose.yml 没透传 TRANSCODE_CONCURRENCY / beat 开关,资源限制与现状不符。
  4. 镜像 97ad0ae2 时期 calculate_asset_quality 因 'str'.value bug 连续失败,8abdeb95 已修复但历史素材 quality_score 仍为 NULL,需要一次性补跑。

变更

(a) 任务路由补全(核心修复)packages/shared/celery_queues.py

所有实时链路任务显式路由到 generation 队列:

任务 队列 说明
worker.generate_video generation 视频生成(已有)
worker.process_tts_synthesis / worker.process_tts_segment_synthesis generation TTS 合成(新增路由)
worker.process_voice_clone generation 音色克隆(新增路由)
worker.extract_voice / worker.extract_background generation 人声/背景提取(新增路由)
ai_avatar_render.execute generation AI 数字人(新增路由)
lipsync_gpu_process_async / lipsync_tts.synthesize_and_submit / lipsync_tts.poll_mediakit_status / lipsync_tts.persist_output_video generation GPU MuseTalk lipsync 全链路(新增路由)
worker.ingest_asset transcode 素材入库(已有)
worker.classify_asset transcode AI 分类(已有)
worker.calculate_asset_quality transcode 素材质量评分(新增路由)
worker.generate_atom_clips / worker.tag_atom_clip / worker.backfill_atom_clip_tags transcode 原子切片 + AI 打标(新增路由)
worker.process_duplication_check / worker.check_duplicate transcode 素材查重(已有,补全 check_duplicate)
worker.batch_download_videos / worker.batch_generate_thumbnails transcode 批量下载/缩略图(新增路由)
worker.cleanup_stale_* (4 个) celery beat 定时巡检/清理(新增显式路由)

(b) Transcode 并发独立可配 infra/docker/entrypoint-worker.sh

  • TRANSCODE_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

  • 容器内三个独立进程:beat(只发定时任务,不消费)+ generation worker(-Q generation)+ transcode worker(-Q transcode,celery)
  • generation worker 不再带 -B,beat 不再与 generation 生命周期耦合
  • 新增 BEAT_ENABLED=1/0 开关,未来独立 beat 容器部署时 worker 容器设为 0 即可
  • 任一进程退出,容器整体退出由 docker restart 拉起

(d) 资源限制与 compose infra/docker/compose.yml

  • 透传 GENERATION_CONCURRENCY/TRANSCODE_CONCURRENCY/BEAT_ENABLED 三个环境变量
  • 健康检查从"任意 celery 进程"改为"celery worker 进程"(beat 本身不算存活依据)
  • 更新注释:当前 2+2 并发场景建议 4C8G;注释写明后续可拆为三个 service 独立扩容
  • .env.example 同步补全三个新变量的说明

(e) 历史 quality 补打分脚本 apps/worker/scripts/backfill_asset_quality.py

一次性脚本(不常驻注册到 celery imports),在 worker 容器内执行:

# 干跑看有多少条要重跑
python -m scripts.backfill_asset_quality --dry-run
# 正式执行(默认分批 50 条/批,每批 sleep 1s,投递到 transcode 队列)
python -m scripts.backfill_asset_quality
# 只重跑最近 30 天
python -m scripts.backfill_asset_quality --since-days 30

不做的事

  • 没把 beat/gen/trans 拆成三个独立 compose service:会连带改 deploy 脚本、watchtower 镜像名、CI pipeline、健康检查。当前单容器三进程已实现"beat 不占槽位 + 独立并发 + 资源限制"全部需求,后续如需独立扩容再拆。
  • 没改 task_enqueue.py 的 WORKER_CONCURRENCY = 4 限流阈值常量(这是业务限流估算用,不影响实际 celery 并发;留待后续与前端排队估算一并调)。
  • 没改 worker.Dockerfile:三进程架构由 entrypoint 实现,镜像层不动。

部署步骤(staging)

  1. PR 合入 develop 后,watchtower 自动拉新 worker 镜像
  2. 在 .env 里按需显式设置(不设则用默认 2/2/1):
    GENERATION_CONCURRENCY=2
    TRANSCODE_CONCURRENCY=2
    BEAT_ENABLED=1
    
  3. 容器启动后验证:
    • 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 在跑
  4. 部署稳定后,在容器里跑一次性脚本:
    docker exec xiaoxia-worker-staging python -m scripts.backfill_asset_quality --dry-run
    docker exec xiaoxia-worker-staging python -m scripts.backfill_asset_quality
    
  5. staging .env 里运维热调的 GENERATION_CONCURRENCY=2 保留,新增显式 TRANSCODE_CONCURRENCY=2。
## 背景 Staging 出现「用户没发多任务却排队」现象。排查发现: 1. `packages/shared/celery_queues.py` 的 `task_routes` 只显式路由了 5 个任务,大量实时任务(TTS、音色克隆、人声/背景提取、AI 数字人、GPU MuseTalk lipsync 全链路)没配路由,默认落到 `celery` 队列被 transcode worker 消费;generation worker 根本收不到这些任务 → 实时任务被后台队列阻塞。 2. `entrypoint-worker.sh` 里 celery beat 嵌在 generation worker 里(`-B`),beat 进程本身不占执行槽但日志/生命周期耦合,且启动并发逻辑中 TRANSCODE_CONCURRENCY 依赖 `WORKER_CONCURRENCY - GENERATION_CONCURRENCY` 差值(WORKER_CONCURRENCY=2, GENERATION_CONCURRENCY=2 时 trans=0→硬编码兜底 1,运维无法独立调)。 3. compose.yml 没透传 `TRANSCODE_CONCURRENCY` / beat 开关,资源限制与现状不符。 4. 镜像 97ad0ae2 时期 `calculate_asset_quality` 因 `'str'.value` bug 连续失败,8abdeb95 已修复但历史素材 quality_score 仍为 NULL,需要一次性补跑。 ## 变更 ### (a) 任务路由补全(核心修复)`packages/shared/celery_queues.py` 所有实时链路任务显式路由到 `generation` 队列: | 任务 | 队列 | 说明 | |---|---|---| | `worker.generate_video` | generation | 视频生成(已有) | | `worker.process_tts_synthesis` / `worker.process_tts_segment_synthesis` | **generation** | TTS 合成(新增路由) | | `worker.process_voice_clone` | **generation** | 音色克隆(新增路由) | | `worker.extract_voice` / `worker.extract_background` | **generation** | 人声/背景提取(新增路由) | | `ai_avatar_render.execute` | **generation** | AI 数字人(新增路由) | | `lipsync_gpu_process_async` / `lipsync_tts.synthesize_and_submit` / `lipsync_tts.poll_mediakit_status` / `lipsync_tts.persist_output_video` | **generation** | GPU MuseTalk lipsync 全链路(新增路由) | | `worker.ingest_asset` | transcode | 素材入库(已有) | | `worker.classify_asset` | transcode | AI 分类(已有) | | `worker.calculate_asset_quality` | **transcode** | 素材质量评分(新增路由) | | `worker.generate_atom_clips` / `worker.tag_atom_clip` / `worker.backfill_atom_clip_tags` | **transcode** | 原子切片 + AI 打标(新增路由) | | `worker.process_duplication_check` / `worker.check_duplicate` | transcode | 素材查重(已有,补全 check_duplicate) | | `worker.batch_download_videos` / `worker.batch_generate_thumbnails` | **transcode** | 批量下载/缩略图(新增路由) | | `worker.cleanup_stale_*` (4 个) | **celery** | beat 定时巡检/清理(新增显式路由) | ### (b) Transcode 并发独立可配 `infra/docker/entrypoint-worker.sh` - `TRANSCODE_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` - 容器内三个独立进程:**beat**(只发定时任务,不消费)+ **generation worker**(`-Q generation`)+ **transcode worker**(`-Q transcode,celery`) - generation worker 不再带 `-B`,beat 不再与 generation 生命周期耦合 - 新增 `BEAT_ENABLED=1/0` 开关,未来独立 beat 容器部署时 worker 容器设为 0 即可 - 任一进程退出,容器整体退出由 docker restart 拉起 ### (d) 资源限制与 compose `infra/docker/compose.yml` - 透传 `GENERATION_CONCURRENCY`/`TRANSCODE_CONCURRENCY`/`BEAT_ENABLED` 三个环境变量 - 健康检查从"任意 celery 进程"改为"celery worker 进程"(beat 本身不算存活依据) - 更新注释:当前 2+2 并发场景建议 4C8G;注释写明后续可拆为三个 service 独立扩容 - `.env.example` 同步补全三个新变量的说明 ### (e) 历史 quality 补打分脚本 `apps/worker/scripts/backfill_asset_quality.py` 一次性脚本(不常驻注册到 celery imports),在 worker 容器内执行: ```bash # 干跑看有多少条要重跑 python -m scripts.backfill_asset_quality --dry-run # 正式执行(默认分批 50 条/批,每批 sleep 1s,投递到 transcode 队列) python -m scripts.backfill_asset_quality # 只重跑最近 30 天 python -m scripts.backfill_asset_quality --since-days 30 ``` ## 不做的事 - 没把 beat/gen/trans 拆成三个独立 compose service:会连带改 deploy 脚本、watchtower 镜像名、CI pipeline、健康检查。当前单容器三进程已实现"beat 不占槽位 + 独立并发 + 资源限制"全部需求,后续如需独立扩容再拆。 - 没改 `task_enqueue.py` 的 `WORKER_CONCURRENCY = 4` 限流阈值常量(这是业务限流估算用,不影响实际 celery 并发;留待后续与前端排队估算一并调)。 - 没改 worker.Dockerfile:三进程架构由 entrypoint 实现,镜像层不动。 ## 部署步骤(staging) 1. PR 合入 develop 后,watchtower 自动拉新 worker 镜像 2. 在 `.env` 里按需显式设置(不设则用默认 2/2/1): ``` GENERATION_CONCURRENCY=2 TRANSCODE_CONCURRENCY=2 BEAT_ENABLED=1 ``` 3. 容器启动后验证: - `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 在跑 4. 部署稳定后,在容器里跑一次性脚本: ```bash docker exec xiaoxia-worker-staging python -m scripts.backfill_asset_quality --dry-run docker exec xiaoxia-worker-staging python -m scripts.backfill_asset_quality ``` 5. staging `.env` 里运维热调的 `GENERATION_CONCURRENCY=2` 保留,新增显式 `TRANSCODE_CONCURRENCY=2`。
xiaoxia added 6 commits 2026-09-28 00:50:49 +08:00
feat(scripts): 一次性脚本 — 对历史 quality_score 缺失的视频素材重新打分 (#2073) - 镜像 97ad0ae2 时期 calculate_asset_quality 因 str.value bug 失败,8abdeb95 已修复; 本脚本扫描 quality_score IS NULL 的视频素材投递到 transcode 队列,支持 dry-run/时间范围/分批限流
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
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 / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (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 / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 48s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m36s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 2m22s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m4s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 3m2s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 3m19s
AI Code Review / AI Code Review (pull_request) Successful in 6m29s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 6m31s
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Style (pull_request) Has been cancelled
CI/CD Pipeline / Unit Tests (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
86086a994f

🚀 预览环境已部署

项目 详情
PR号 #2079
预览链接 https://pr-2079.preview.xiaoxiajianji.com
API环境 staging

💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。

🔄 每次提交新代码后预览环境会自动更新。

🗑️ PR 关闭或合并后,预览环境会自动清理。

🚀 **预览环境已部署** | 项目 | 详情 | |------|------| | PR号 | #2079 | | 预览链接 | [https://pr-2079.preview.xiaoxiajianji.com](https://pr-2079.preview.xiaoxiajianji.com) | | API环境 | staging | > 💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。 > > 🔄 每次提交新代码后预览环境会自动更新。 > > 🗑️ PR 关闭或合并后,预览环境会自动清理。
auto-approve-bot added 1 commit 2026-09-28 00:57:26 +08:00
style: auto-format with black + isort + ruff + prettier [skip ci-format-check]
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (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 / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 50s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 54s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m18s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m11s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m55s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m11s
AI Code Review / AI Code Review (pull_request) Successful in 6m48s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 7m5s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 9m23s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 11m17s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 1s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 8m52s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 32s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 45s
efba58f1bd
auto-approve-bot merged commit 01991f14d7 into develop 2026-09-28 01:09:31 +08:00
auto-approve-bot deleted branch fix/worker-queue-split 2026-09-28 01:09:33 +08:00

🗑️ 预览环境已清理

PR #2079 已关闭或合并,对应的预览环境已被清理。

如有需要,可以重新打开 PR 来重新生成预览环境。

🗑️ **预览环境已清理** PR #2079 已关闭或合并,对应的预览环境已被清理。 > 如有需要,可以重新打开 PR 来重新生成预览环境。
Sign in to join this conversation.