Files
xiaoxia-saas/packages/shared/celery_queues.py
T
xiaoxia 01991f14d7
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 3s
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 3s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Check push changed paths (push) Successful in 13s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging 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 Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m3s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m7s
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 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 / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m9s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 3m13s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m24s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 2m17s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m33s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Integration Tests (push) Successful in 4m17s
CI/CD Pipeline / Validate - Style (push) Successful in 4m51s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 4m53s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 55s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m42s
AI Code Review / AI Code Review (pull_request) Successful in 6m57s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m52s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m45s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 9m5s
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 4m58s
CI/CD Pipeline / Unit Tests (push) Successful in 11m33s
CI/CD Pipeline / Validate - Security (push) Successful in 11m53s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
fix(queue): #2073 Worker 队列分流——任务路由补全 + beat 独立 + transcode 并发独立 (#2079)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-28 01:09:29 +08:00

89 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