diff --git a/packages/shared/celery_queues.py b/packages/shared/celery_queues.py index c9f4e45e6..ffb975746 100644 --- a/packages/shared/celery_queues.py +++ b/packages/shared/celery_queues.py @@ -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,47 @@ 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 +80,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