Files
xiaoxia-saas/packages/shared/celery_queues.py
T

91 lines
4.9 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Celery 队列定义与路由配置(API / Worker 共享)。
#1714 + #2073 队列分流:用户同步等待的实时任务路由到 `generation` 高优队列,
由专用 generation worker 独占消费;素材入库/转码/AI 分析/查重等后台批量任务路由
到 `transcode` 队列;beat 定时清理等轻量维护任务走默认 `celery` 队列。
transcode / celery 队列积压时,generation 队列仍能被立即领取,不阻塞用户实时链路。
队列说明:
- generation: 用户同步等待的实时任务(视频生成、TTS、音色克隆、lipsync、AI 数字人、人声/背景提取)
- transcode: 后台批量/异步任务(素材入库转码、AI 分类打标、质量评分、原子切片、查重、批量下载/缩略图)
- celery: beat 定时巡检/清理等轻量维护任务(极短、低优、不占业务槽)
"""
from __future__ import annotations
from kombu import Queue
# ── 队列名常量(生产端与消费端共用,禁止拼写漂移) ──
QUEUE_GENERATION = "generation"
QUEUE_TRANSCODE = "transcode"
QUEUE_DEFAULT = "celery"
# 三个消费组各自消费的队列列表(顺序即优先级:高优队列排在前面)
WORKER_QUEUES_GENERATION = (QUEUE_GENERATION,)
WORKER_QUEUES_TRANSCODE = (QUEUE_TRANSCODE, QUEUE_DEFAULT)
# 队列声明:持久化队列,broker 重启不丢消息
task_queues = (
Queue(QUEUE_GENERATION, routing_key=QUEUE_GENERATION, durable=True),
Queue(QUEUE_TRANSCODE, routing_key=QUEUE_TRANSCODE, durable=True),
Queue(QUEUE_DEFAULT, routing_key=QUEUE_DEFAULT, durable=True),
)
# ── 任务路由表:task name → 队列 ──
# 键支持 celery 标准通配符。所有生产端(API send_task / worker 内 send_task)
# 未显式指定 queue 时按此表路由;漏配会走默认队列 celery,被 transcode worker 消费。
# 新增实时任务务必在此表显式路由到 generation,避免落到后台队列排队。
task_routes = {
# ── 高优先级:用户同步等待的实时链路 ──
# 视频生成(主链路)
"worker.generate_video": {"queue": QUEUE_GENERATION},
# 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 占住不调度
GENERATION_WORKER_PREFETCH_MULTIPLIER = 1
def apply_queue_settings(app) -> None:
"""把队列分流配置应用到 Celery app(API 生产端与 Worker 消费端都要调用)。
配置 task_queues / task_routes / task_default_queue。生产端靠 task_routes
把消息投递到对应队列;消费端靠启动参数 -Q 控制自己消费哪些队列(entrypoint)。
"""
app.conf.task_queues = task_queues
app.conf.task_routes = task_routes
app.conf.task_default_queue = QUEUE_DEFAULT