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
6 changed files with 289 additions and 118 deletions
+21 -2
View File
@@ -79,14 +79,33 @@ CELERY_BROKER_URL=redis://localhost:6379/0
CELERY_RESULT_BACKEND=redis://localhost:6379/1
# ==================== Worker 配置 ====================
# ==================== Worker 配置(#2073 队列分流) ====================
#
# 容器内跑三个独立进程:beat(只发定时任务)+ generation worker(实时高优)
# + transcode worker(后台批量/清理)。三个进程的并发与开关独立配置。
# Worker 进程名称
WORKER_NAME=xiaoxia-saas-worker
# Worker 并发数(同时执行的任务数)
# 总并发参考(兼容旧变量):
# - 若 GENERATION_CONCURRENCY 与 TRANSCODE_CONCURRENCY 都未显式设置,
# entrypoint 会按此总数对半分配(gen=ceil(total/2), trans=剩余,各至少 1);
# - 任一个 *_CONCURRENCY 显式设置后,按显式值生效,忽略此变量对应部分。
WORKER_CONCURRENCY=4
# Generation worker 并发数(用户实时任务:视频生成/TTS/音色克隆/lipsync/数字人)
# 实时链路对延迟敏感,建议 2C 以上机器设为 2;高负载场景可加到 4。
GENERATION_CONCURRENCY=2
# Transcode worker 并发数(后台批量:素材入库转码/AI 分类打标/质量评分/查重/批量下载)
# 后台任务可排队,独立伸缩;素材入库量大时可加到 4。
TRANSCODE_CONCURRENCY=2
# 是否在本容器启动 celery beat 进程(默认 1)。
# 默认 beat 与 worker 同容器部署;若要独立 beat 容器部署,worker 容器设为 0、
# beat 容器单独跑 `celery -A worker_app.celery_app beat` 并设 BEAT_ENABLED=1。
BEAT_ENABLED=1
# 每个子进程最多处理多少任务后重启(防止内存泄漏)
WORKER_MAX_TASKS_PER_CHILD=1000
View File
@@ -0,0 +1,112 @@
"""一次性脚本:对历史 quality_score 缺失的视频素材重新打分。
背景(#2073):镜像 97ad0ae2 时期 calculate_quality_score / classify_from_analysis
返回 str 而非 AssetClassification 枚举,导致 calculate_asset_quality 连续报
"'str' object has no attribute 'value'",大量视频素材的 quality_score 卡在 NULL。
镜像 8abdeb95 已修复枚举 bug,但历史失败记录不会自动重跑。本脚本扫描全表,
把 quality_score IS NULL 的视频素材重新投递到 worker.calculate_asset_quality 任务。
使用方式(在 worker 容器内执行):
cd /app/apps/worker
# 干跑,只打印会重跑多少条,不发任务
python -m scripts.backfill_asset_quality --dry-run
# 正式执行
python -m scripts.backfill_asset_quality
# 只重跑最近 N 天的
python -m scripts.backfill_asset_quality --since-days 30
# 限流:每投递一批 sleep 几秒,避免瞬间打爆 transcode 队列
python -m scripts.backfill_asset_quality --batch-size 50 --sleep 2
也可以直接在 staging 机器上 exec 进容器:
docker exec -e PYTHONPATH=/app:/app/apps/api:/app/packages xiaoxia-worker-staging \
python -m scripts.backfill_asset_quality --dry-run
"""
from __future__ import annotations
import argparse
# 保证可以以 python -m scripts.xxx 在容器 /app/apps/worker 下执行
# 也兼容在 repo 根目录下执行(注入路径)
import os
import sys
import time
from datetime import UTC, datetime, timedelta
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
_WORKER_DIR = os.path.dirname(_SCRIPT_DIR) # apps/worker
_APPS_DIR = os.path.dirname(_WORKER_DIR) # apps
_REPO_ROOT = os.path.dirname(_APPS_DIR) # repo root
for p in (_REPO_ROOT, os.path.join(_REPO_ROOT, "apps", "api"), _REPO_ROOT):
if p not in sys.path:
sys.path.insert(0, p)
def main() -> int:
parser = argparse.ArgumentParser(description="补打历史视频素材 quality_score")
parser.add_argument("--dry-run", action="store_true", help="只统计数量,不投递任务")
parser.add_argument("--since-days", type=int, default=0, help="只处理最近 N 天上传的素材(0=全部)")
parser.add_argument("--batch-size", type=int, default=50, help="每批投递数量,默认 50")
parser.add_argument("--sleep", type=float, default=1.0, help="批次之间 sleep 秒数,默认 1s")
parser.add_argument("--queue", type=str, default="transcode", help="投递队列(默认 transcode)")
args = parser.parse_args()
# 延迟 import,避免在 dry-run 时依赖完整 DB 环境
from worker_app.celery_app import celery_app
from worker_app.db import SessionLocal
from packages.adapters.sqlalchemy_impl.models import AssetModel
db = SessionLocal()
try:
q = db.query(AssetModel).filter(
AssetModel.file_type == "video",
AssetModel.quality_score.is_(None),
)
if args.since_days > 0:
cutoff = datetime.now(UTC) - timedelta(days=args.since_days)
q = q.filter(AssetModel.created_at >= cutoff)
# 先 count 打印
total = q.count()
print(
f"[backfill] 待重跑 quality_score 的视频素材: {total} 条"
f"{' (dry-run,不投递)' if args.dry_run else ''}"
f"{' (最近 ' + str(args.since_days) + ' 天)' if args.since_days > 0 else ''}",
flush=True,
)
if total == 0 or args.dry_run:
return 0
# 分批投递
submitted = 0
batch = 0
offset = 0
while True:
assets = q.order_by(AssetModel.created_at.desc()).offset(offset).limit(args.batch_size).all()
if not assets:
break
batch += 1
for a in assets:
try:
celery_app.send_task(
"worker.calculate_asset_quality",
args=[a.id],
queue=args.queue,
)
submitted += 1
except Exception as e: # noqa: BLE001
print(f"[backfill] 投递失败 asset_id={a.id}: {e}", flush=True)
print(f"[backfill] batch {batch}: 已累计投递 {submitted}/{total}", flush=True)
offset += len(assets)
if args.sleep > 0 and offset < total:
time.sleep(args.sleep)
print(f"[backfill] 完成,共投递 {submitted} 条任务到 {args.queue} 队列", flush=True)
return 0
finally:
db.close()
if __name__ == "__main__":
sys.exit(main())
+44 -72
View File
@@ -10,20 +10,25 @@
# API_IMAGE - API 镜像名称 (默认: xiaoxia-saas-api:dev)
# WORKER_IMAGE - Worker 镜像名称 (默认: xiaoxia-saas-worker:dev)
# WEB_IMAGE - Web 镜像名称 (默认: xiaoxia-saas-web:dev)
# WEB_DOCKERFILE - Web Dockerfile 路径
# WEB_NGINX_CONF - Nginx 配置文件路径
# API_PORT - API 端口映射 (staging: 8000, production: 8001)
# WEB_PORT - Web 端口映射 (staging: 3001, production: 3002)
# GENERATED_FILES_HOST_DIR - 生成文件的主机目录
# WORKER_CONCURRENCY - Worker 并发数 (默认: 4)
# GENERATION_CONCURRENCY - Generation worker 并发(用户实时任务,默认 2)
# TRANSCODE_CONCURRENCY - Transcode worker 并发(后台/转码/AI,默认 2)
# WORKER_MAX_TASKS_PER_CHILD - Worker 每个子进程最大任务数 (默认: 100)
# BEAT_ENABLED - 容器内启动 celery beat(默认 1;独立 beat 容器部署设为 0)
# WORKER_CONCURRENCY - 兼容旧变量:未显式设置上面两个并发时按此总数分配
#
# 重要:
# 重要:
# - 生产环境不要挂载 web-dist volume,否则会导致 403
# - 确保环境隔离网络已创建: docker network create xiaoxia-net-${ENV}
# - ENV=staging → xiaoxia-net-staging
# - ENV=production → xiaoxia-net-production
# - ENV=staging -> xiaoxia-net-staging
# - ENV=production -> xiaoxia-net-production
#
# #2073 队列分流:worker 容器内跑三个独立进程——beat(只发定时任务)、
# generation worker(只消费 generation 队列,实时高优)、transcode worker(消费
# transcode + celery 队列,后台任务)。beat 不再嵌入 generation worker,
# 不占实时任务槽位;TRANSCODE_CONCURRENCY 独立伸缩,不再依赖 WORKER_CONCURRENCY 差值。
# ===========================================
# 日志轮转配置(所有服务共享)
@@ -40,53 +45,39 @@ services:
# =========================================
api:
image: ${API_IMAGE:-xiaoxia-saas-api:dev}
# 不在生产环境构建镜像,使用预构建的镜像
# build:
# context: ../..
# dockerfile: infra/docker/api.Dockerfile
container_name: xiaoxia-api-${ENV:-staging}
restart: unless-stopped
stop_grace_period: 30s
stop_signal: SIGTERM
# 环境变量文件(包含数据库密码等敏感信息)
env_file:
- ../../.env
environment:
APP_ENV: ${APP_ENV:-staging}
GENERATED_FILES_DIR: /app/generated
GENERATED_FILES_URL_PREFIX: /generated-files
PUBLIC_API_BASE_URL: ${PUBLIC_API_BASE_URL:-https://api.xiaoxiajianji.com}
# 端口映射
# Staging: 8000 -> 8000
# Production: 8001 -> 8000
ports:
- "127.0.0.1:${API_PORT:-8000}:8000"
# 共享生成文件目录 + 抖音 cookies 等运行时配置
volumes:
- generated-files:/app/generated
- ../../deploy/configs:/app/configs:ro
networks:
- xiaoxia-net
# 健康检查配置
healthcheck:
test: ["CMD", "python", "-c", "import urllib.request; urllib.request.urlopen('http://localhost:8000/health', timeout=5)"]
interval: 30s
timeout: 10s
retries: 3
start_period: 40s
logging: *default-logging
# =========================================
# 资源限制建议(生产环境建议启用)
# =========================================
deploy:
resources:
limits:
@@ -97,39 +88,48 @@ services:
memory: 512M
# =========================================
# Worker 服务(Celery 任务队列)
# Worker 服务(#2073 队列分流:beat + generation + transcode 同容器三进程)
# =========================================
# 三个进程独立启动,任一退出则容器整体退出由 docker restart 拉起;
# 各自的并发与资源占用通过环境变量控制:
# - generation:GENERATION_CONCURRENCY(默认 2),消费 generation 队列
# - transcode: TRANSCODE_CONCURRENCY(默认 2),消费 transcode,celery 队列
# - beat: 不消费任务,只发定时任务到 celery 默认队列
worker:
image: ${WORKER_IMAGE:-xiaoxia-saas-worker:dev}
container_name: xiaoxia-worker-${ENV:-staging}
restart: unless-stopped
# 长任务(ingest HEVC 转码最长 30min、生成硬超时 11min)给足优雅关闭窗口
stop_grace_period: 300s
stop_signal: SIGTERM
env_file:
- ../../.env
environment:
APP_ENV: ${APP_ENV:-staging}
# 兼容旧变量:若两个 *_CONCURRENCY 均未显式设置,entrypoint 会按此总数分配
WORKER_CONCURRENCY: ${WORKER_CONCURRENCY:-4}
WORKER_MAX_TASKS_PER_CHILD: ${WORKER_MAX_TASKS_PER_CHILD:-100}
# #1714 队列隔离:generation 队列独占 worker(默认并发 2),其余并发给转码
# #2073 队列独立伸缩:generation 默认 2,transcode 默认 2(不再差值计算)
GENERATION_CONCURRENCY: ${GENERATION_CONCURRENCY:-2}
TRANSCODE_CONCURRENCY: ${TRANSCODE_CONCURRENCY:-2}
# beat 默认在本容器启动;独立 beat 容器部署时设为 0
BEAT_ENABLED: ${BEAT_ENABLED:-1}
GENERATED_FILES_DIR: /app/generated
GENERATED_FILES_URL_PREFIX: /generated-files
PUBLIC_API_BASE_URL: ${PUBLIC_API_BASE_URL:-https://api.xiaoxiajianji.com}
volumes:
- generated-files:/app/generated
networks:
- xiaoxia-net
# 健康检查配置
# 注:容器内无 pgrep/ps,扫描 /proc 所有进程的 cmdline 查找 celery 进程
# 健康检查:至少有一个 celery worker 进程在跑(beat 本身不作为存活依据)
healthcheck:
test: ["CMD-SHELL", "grep -lq celery /proc/[0-9]*/cmdline 2>/dev/null || exit 1"]
test: ["CMD-SHELL", "grep -q 'celery.*worker' /proc/[0-9]*/cmdline 2>/dev/null || exit 1"]
interval: 30s
timeout: 10s
retries: 3
@@ -137,12 +137,8 @@ services:
logging: *default-logging
# =========================================
# 资源限制建议(生产环境建议启用)
# =========================================
# 注意: Worker 需要处理视频,建议分配更多资源
# #1714 队列隔离后容器内运行 generation + transcode 两个 worker 进程,
# 总并发 = WORKER_CONCURRENCY(默认 4),4C8G 以上确保视频渲染不 OOM
# 资源限制:容器总资源 = gen + trans + beat,按 2+2 并发场景建议 4C8G;
# 后续如需独立扩容/重启,可拆为 worker-generation / worker-transcode / worker-beat 三个 service。
deploy:
resources:
limits:
@@ -157,35 +153,21 @@ services:
# =========================================
web:
image: ${WEB_IMAGE:-xiaoxia-saas-web:dev}
# 不在生产环境构建镜像,使用 web-artifact.Dockerfile
# build:
# context: ../..
# dockerfile: ${WEB_DOCKERFILE:-infra/docker/web.Dockerfile}
# args:
# (NGINX_CONF no longer needed - all configs baked into image)
container_name: xiaoxia-web-${ENV:-staging}
restart: unless-stopped
# 端口映射
# Staging: 3001 -> 80
# Production: 3002 -> 80 (通过 Nginx 反向代理)
ports:
- "127.0.0.1:${WEB_PORT:-3001}:80"
networks:
- xiaoxia-net
# =========================================
# Nginx 配置运行时覆盖(双保险:entrypoint 也按 APP_ENV 选择配置)
# 确保容器使用正确环境的 nginx 配置,即使镜像构建时使用了默认配置
# 注意: 只覆盖 /etc/nginx/conf.d/default.conf,不挂载 /usr/share/nginx/html
# =========================================
environment:
- APP_ENV=${ENV:-staging}
volumes:
- ./nginx-${ENV:-staging}.conf:/etc/nginx/conf.d/default.conf:ro
healthcheck:
test: ["CMD", "wget", "--spider", "-q", "http://127.0.0.1:80"]
interval: 30s
@@ -194,9 +176,6 @@ services:
logging: *default-logging
# =========================================
# 资源限制建议
# =========================================
deploy:
resources:
limits:
@@ -212,9 +191,6 @@ volumes:
driver_opts:
type: none
o: bind
# 重要: 确保主机目录存在且有正确权限
# Staging: /var/lib/xiaoxia-saas-staging/generated
# Production: /var/lib/xiaoxia-saas-production/generated
device: ${GENERATED_FILES_HOST_DIR:?GENERATED_FILES_HOST_DIR must be set in .env}
# ===========================================
@@ -223,8 +199,4 @@ volumes:
networks:
xiaoxia-net:
external: true
# 网络名根据 ENV 变量区分,实现 staging/production 环境隔离
# staging: xiaoxia-net-staging
# production: xiaoxia-net-production
name: xiaoxia-net-${ENV:-staging}
+67 -29
View File
@@ -1,48 +1,77 @@
#!/bin/bash
# Worker 启动脚本 — #1714 队列隔离
# Worker 启动脚本 — #1714 + #2073 队列分流
#
# 部署约束:worker 容器单实例(replicas=1),容器内启动两个 celery 进程:
# 1. generation-worker:独占消费 generation 队列(用户视频生成,高优先级),
# 内嵌 celery beat(-B),定时清理任务只在一个进程里跑,避免重复执行;
# 2. transcode-worker:消费 transcode + celery 默认队列(素材转码/分类/查重/
# 配音/下载等后台任务)。
# 转码队列积压时,generation 队列仍有独立 worker 立即领取视频生成任务。
# 容器内启动三个独立进程(任一退出则整体退出由 docker restart 拉起):
# 1. beat:celery beat 调度器,不消费任何任务,只发定时任务到 celery 默认队列
# 2. generation-worker:独占消费 generation 队列(用户实时任务,高优先级)
# 3. transcode-worker:消费 transcode + celery 默认队列(后台/清理任务)
#
# 环境变量:
# WORKER_CONCURRENCY 总并发槽参考(默认 4);生成 worker 并发默认 2,
# 可用 GENERATION_CONCURRENCY 覆盖
# GENERATION_CONCURRENCY generation worker 并发(默认 2)
# TRANSCODE_CONCURRENCY transcode worker 并发(默认 = WORKER_CONCURRENCY - 2,最小 1)
# TRANSCODE_CONCURRENCY transcode worker 并发(默认 2)
# WORKER_MAX_TASKS_PER_CHILD 每个子进程最大任务数(默认 100)
# WORKER_CONCURRENCY 兼容旧变量:若未显式设置 GENERATION_CONCURRENCY /
# TRANSCODE_CONCURRENCY,则按比例分配(gen=ceil(total*1/2),
# trans=剩余,各至少 1);已显式设置时忽略此变量。
# BEAT_ENABLED 是否在本容器内启动 beat 进程(默认 1);
# 若独立 beat 容器部署设为 0。
set -e
CONCURRENCY="${WORKER_CONCURRENCY:-4}"
MAX_TASKS="${WORKER_MAX_TASKS_PER_CHILD:-100}"
GEN_CONCURRENCY="${GENERATION_CONCURRENCY:-2}"
if [ -z "$TRANSCODE_CONCURRENCY" ]; then
TRANS_CONCURRENCY=$((CONCURRENCY - GEN_CONCURRENCY))
if [ "$TRANS_CONCURRENCY" -lt 1 ]; then
TRANS_CONCURRENCY=1
fi
# ── 并发计算:显式 env 优先;否则从 WORKER_CONCURRENCY 按比例推导 ──
if [ -n "$GENERATION_CONCURRENCY" ]; then
GEN_CONCURRENCY="$GENERATION_CONCURRENCY"
else
TRANS_CONCURRENCY="$TRANSCODE_CONCURRENCY"
TOTAL="${WORKER_CONCURRENCY:-4}"
GEN_CONCURRENCY=$(( (TOTAL + 1) / 2 ))
if [ "$GEN_CONCURRENCY" -lt 1 ]; then GEN_CONCURRENCY=1; fi
fi
echo "Starting generation worker (queue=generation, concurrency=$GEN_CONCURRENCY, beat embedded)"
if [ -n "$TRANSCODE_CONCURRENCY" ]; then
TRANS_CONCURRENCY="$TRANSCODE_CONCURRENCY"
else
if [ -n "$WORKER_CONCURRENCY" ] && [ -z "$GENERATION_CONCURRENCY" ]; then
# 两个都没显式设置,按 WORKER_CONCURRENCY 分配剩余
TOTAL="$WORKER_CONCURRENCY"
TRANS_CONCURRENCY=$(( TOTAL - GEN_CONCURRENCY ))
if [ "$TRANS_CONCURRENCY" -lt 1 ]; then TRANS_CONCURRENCY=1; fi
else
# 默认 2(#2073:独立伸缩,不再依赖 WORKER_CONCURRENCY 差值)
TRANS_CONCURRENCY=2
fi
fi
BEAT_ENABLED="${BEAT_ENABLED:-1}"
PIDS=()
# ── 1. Beat 调度器(独立进程,不消费任务)──
if [ "$BEAT_ENABLED" = "1" ] || [ "$BEAT_ENABLED" = "true" ]; then
echo "Starting beat scheduler (schedule file=/tmp/celerybeat-schedule)"
celery \
-A worker_app.celery_app \
beat \
--loglevel=info \
-s /tmp/celerybeat-schedule &
PIDS+=($!)
fi
# ── 2. Generation worker(实时高优队列)──
echo "Starting generation worker (queue=generation, concurrency=$GEN_CONCURRENCY)"
celery \
-A worker_app.celery_app \
worker \
--loglevel=info \
"-B" \
-s /tmp/celerybeat-schedule \
-Q generation \
"--concurrency=${GEN_CONCURRENCY}" \
"--max-tasks-per-child=${MAX_TASKS}" \
-n generation@%h &
GEN_PID=$!
PIDS+=($!)
GEN_PID=${PIDS[1]:-${PIDS[0]}}
# ── 3. Transcode worker(后台 + 清理队列)──
echo "Starting transcode worker (queues=transcode,celery, concurrency=$TRANS_CONCURRENCY)"
celery \
-A worker_app.celery_app \
@@ -52,13 +81,22 @@ celery \
"--concurrency=${TRANS_CONCURRENCY}" \
"--max-tasks-per-child=${MAX_TASKS}" \
-n transcode@%h &
TRANS_PID=$!
PIDS+=($!)
TRANS_PID=${PIDS[2]:-${PIDS[1]}}
# 任一进程退出则终止另一个,让容器整体重启(restart: unless-stopped)
trap 'echo "Shutting down workers..."; kill -TERM $GEN_PID $TRANS_PID 2>/dev/null || true' TERM INT
# 任一进程退出则终止其他进程,让容器整体重启
cleanup() {
echo "Shutting down all celery processes..."
for pid in "${PIDS[@]}"; do
kill -TERM "$pid" 2>/dev/null || true
done
}
trap cleanup TERM INT
wait -n $GEN_PID $TRANS_PID
# wait -n 等待任意一个子进程退出(bash 4.3+)
# 容器镜像基础为 python:3.11-slim,bash 版本满足
wait -n "${PIDS[@]}"
EXIT_CODE=$?
echo "One worker exited (code=$EXIT_CODE), stopping the other..."
kill -TERM $GEN_PID $TRANS_PID 2>/dev/null || true
exit $EXIT_CODE
echo "One celery process exited (code=$EXIT_CODE), stopping the rest..."
cleanup
exit "$EXIT_CODE"
+45 -15
View File
@@ -1,14 +1,14 @@
"""Celery 队列定义与路由配置(API / Worker 共享)。
#1714 队列隔离:用户等待的视频生成任务路由到高优先级 `generation` 队列,
由专用 worker 进程独占消费;素材入库/转码等后台批量任务路由到 `transcode`
队列;其余杂项任务走默认 `celery` 队列。转码队列积压时,视频生成任务
仍能被 generation worker 立即领取执行,不会排队。
#1714 + #2073 队列分流:用户同步等待的实时任务路由到 `generation` 高优队列,
由专用 generation worker 独占消费;素材入库/转码/AI 分析/查重等后台批量任务路由
到 `transcode` 队列;beat 定时清理等轻量维护任务走默认 `celery` 队列。
transcode / celery 队列积压时,generation 队列仍能被立即领取,不阻塞用户实时链路。
队列说明:
- generation: 用户提交的视频生成/预览渲染(延迟敏感,资源消耗大)
- transcode: 素材入库(HEVC 转码)、AI 分类、素材查重(批量、可排队)
- celery(默认): 配音、语音、下载缩略图、定时清理等杂项
- generation: 用户同步等待的实时任务(视频生成、TTS、音色克隆、lipsync、AI 数字人、人声/背景提取)
- transcode: 后台批量/异步任务(素材入库转码、AI 分类打标、质量评分、原子切片、查重、批量下载/缩略图)
- celery: beat 定时巡检/清理等轻量维护任务(极短、低优、不占业务槽)
"""
from __future__ import annotations
@@ -20,8 +20,9 @@ QUEUE_GENERATION = "generation"
QUEUE_TRANSCODE = "transcode"
QUEUE_DEFAULT = "celery"
# Worker 消费的队列列表(顺序即优先级:高优队列排在前面)
WORKER_QUEUES = (QUEUE_GENERATION, QUEUE_TRANSCODE, QUEUE_DEFAULT)
# 三个消费组各自消费的队列列表(顺序即优先级:高优队列排在前面)
WORKER_QUEUES_GENERATION = (QUEUE_GENERATION,)
WORKER_QUEUES_TRANSCODE = (QUEUE_TRANSCODE, QUEUE_DEFAULT)
# 队列声明:持久化队列,broker 重启不丢消息
task_queues = (
@@ -31,15 +32,45 @@ task_queues = (
)
# ── 任务路由表:task name → 队列 ──
# 键支持 celery 标准通配符。
# 键支持 celery 标准通配符。所有生产端(API send_task / worker 内 send_task)
# 未显式指定 queue 时按此表路由;漏配会走默认队列 celery,被 transcode worker 消费。
# 新增实时任务务必在此表显式路由到 generation,避免落到后台队列排队。
task_routes = {
# 高优先级:用户等待的视频生成
# ── 高优先级:用户同步等待的实时链路 ──
# 视频生成(主链路)
"worker.generate_video": {"queue": QUEUE_GENERATION},
# 后台批量:素材入库/转码 + AI 分类 + 素材查重,积压不影响生成
# TTS 合成 / 片段合成(配音页、视频生成配乐/TTS 链路)
"worker.process_tts_synthesis": {"queue": QUEUE_GENERATION},
"worker.process_tts_segment_synthesis": {"queue": QUEUE_GENERATION},
# 音色克隆(用户主动上传样本等待克隆完成)
"worker.process_voice_clone": {"queue": QUEUE_GENERATION},
# 人声/背景提取(音色克隆前置步骤,用户同步等待)
"worker.extract_voice": {"queue": QUEUE_GENERATION},
"worker.extract_background": {"queue": QUEUE_GENERATION},
# AI 数字人渲染(用户主动触发,等待成片)
"ai_avatar_render.execute": {"queue": QUEUE_GENERATION},
# GPU MuseTalk 口型同步(用户等成片,链路子任务全部走 generation 避免跨队列阻塞)
"lipsync_gpu_process_async": {"queue": QUEUE_GENERATION},
"lipsync_tts.synthesize_and_submit": {"queue": QUEUE_GENERATION},
"lipsync_tts.poll_mediakit_status": {"queue": QUEUE_GENERATION},
"lipsync_tts.persist_output_video": {"queue": QUEUE_GENERATION},
# ── 后台批量:素材入库/转码 + AI 分析/打标 + 查重,积压不影响生成 ──
"worker.ingest_asset": {"queue": QUEUE_TRANSCODE},
"worker.classify_asset": {"queue": QUEUE_TRANSCODE},
"worker.calculate_asset_quality": {"queue": QUEUE_TRANSCODE},
"worker.generate_atom_clips": {"queue": QUEUE_TRANSCODE},
"worker.tag_atom_clip": {"queue": QUEUE_TRANSCODE},
"worker.backfill_atom_clip_tags": {"queue": QUEUE_TRANSCODE},
"worker.process_duplication_check": {"queue": QUEUE_TRANSCODE},
"worker.check_duplicate": {"queue": QUEUE_TRANSCODE},
"worker.batch_download_videos": {"queue": QUEUE_TRANSCODE},
"worker.batch_generate_thumbnails": {"queue": QUEUE_TRANSCODE},
# ── beat 定时清理/巡检任务走默认 celery 队列(由 transcode worker 消费)──
# 未在此表显式列出的 cleanup 任务会落到默认队列 celery,不占 generation 槽位。
"worker.cleanup_stale_pending_tasks": {"queue": QUEUE_DEFAULT},
"worker.cleanup_stale_running_tasks": {"queue": QUEUE_DEFAULT},
"worker.cleanup_stale_ingest_jobs": {"queue": QUEUE_DEFAULT},
"worker.cleanup_stale_voice_clones": {"queue": QUEUE_DEFAULT},
}
# 生成任务的预取数:渲染是长任务,预取 1 避免任务被某个 worker 占住不调度
@@ -47,11 +78,10 @@ GENERATION_WORKER_PREFETCH_MULTIPLIER = 1
def apply_queue_settings(app) -> None:
"""把队列隔离配置应用到 Celery app(API 生产端与 Worker 消费端都要调用)。
"""把队列分流配置应用到 Celery app(API 生产端与 Worker 消费端都要调用)。
配置 task_queues / task_routes / task_default_queue。生产端靠 task_routes
把消息投递到对应队列;消费端靠 task_queues 声明自己消费哪些队列
(实际消费集由启动参数 -Q 控制)。
把消息投递到对应队列;消费端靠启动参数 -Q 控制自己消费哪些队列(entrypoint)。
"""
app.conf.task_queues = task_queues
app.conf.task_routes = task_routes