From 54deafd7b08322a4df15c8d74851efaa2f5188ca Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 28 Sep 2026 00:49:48 +0800 Subject: [PATCH] =?UTF-8?q?fix(queue):=20=E8=A1=A5=E5=85=A8=20task=5Froute?= =?UTF-8?q?s=EF=BC=8C=E6=89=80=E6=9C=89=E5=AE=9E=E6=97=B6=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E6=98=BE=E5=BC=8F=E8=B7=AF=E7=94=B1=E5=88=B0=20genera?= =?UTF-8?q?tion=20=E9=98=9F=E5=88=97=20(#2073)=20-=20TTS/=E9=9F=B3?= =?UTF-8?q?=E8=89=B2=E5=85=8B=E9=9A=86/lipsync/AI=E6=95=B0=E5=AD=97?= =?UTF-8?q?=E4=BA=BA/=E4=BA=BA=E5=A3=B0=E6=8F=90=E5=8F=96=E7=AD=89?= =?UTF-8?q?=E5=AE=9E=E6=97=B6=E9=93=BE=E8=B7=AF=E5=85=A8=E9=83=A8=E8=B5=B0?= =?UTF-8?q?=20generation=EF=BC=8C=E5=90=8E=E5=8F=B0=E5=8E=9F=E5=AD=90?= =?UTF-8?q?=E5=88=87=E7=89=87/=E6=89=93=E6=A0=87/=E6=89=B9=E9=87=8F?= =?UTF-8?q?=E4=B8=8B=E8=BD=BD/=E7=BC=A9=E7=95=A5=E5=9B=BE=E8=B5=B0=20trans?= =?UTF-8?q?code=EF=BC=8Cbeat=20=E6=B8=85=E7=90=86=E8=B5=B0=20celery=20?= =?UTF-8?q?=E9=BB=98=E8=AE=A4=E9=98=9F=E5=88=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/shared/celery_queues.py | 62 ++++++++++++++++++++++++-------- 1 file changed, 47 insertions(+), 15 deletions(-) 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