from celery import Celery from worker_app.core.config import get_settings settings = get_settings() celery_app = Celery(settings.worker_name) celery_app.conf.broker_url = settings.broker_url celery_app.conf.result_backend = settings.result_backend celery_app.conf.broker_connection_retry_on_startup = True # #1714 队列隔离:generation(高优,独占 worker)/ transcode(素材转码)/ celery(默认) from packages.shared.celery_queues import ( # noqa: E402 GENERATION_WORKER_PREFETCH_MULTIPLIER, apply_queue_settings, ) apply_queue_settings(celery_app) # 长渲染任务预取 1,避免任务被预取占住导致调度不均 celery_app.conf.worker_prefetch_multiplier = GENERATION_WORKER_PREFETCH_MULTIPLIER celery_app.conf.task_acks_late = True # worker 崩溃时未完成任务重回队列,由执行前守卫丢弃作废消息 celery_app.conf.imports = ( "worker_app.tasks.health", "worker_app.tasks.ingest", "worker_app.tasks.classification", "worker_app.tasks.generation", "worker_app.tasks.voice_extraction", "worker_app.tasks.voice_clone", "worker_app.tasks.tts_synthesis", "worker_app.tasks.batch_download", "worker_app.tasks.duplication_check", "worker_app.tasks._startup", "apps.worker.video_processing.dedup", "worker_app.tasks.cleanup", ) # Celery Beat 定时任务调度 # 注:worker 单实例内嵌 beat(entrypoint-worker.sh -B),定时任务不会重复执行 celery_app.conf.beat_schedule = { # pending 任务超时清理:worker 停止消费后,卡 pending 的任务 15 分钟内释放限流名额 "cleanup-stale-pending-tasks": { "task": "worker.cleanup_stale_pending_tasks", "schedule": 300.0, # 每 5 分钟(秒) "options": {"expires": 240}, # 4 分钟过期,避免堆积 }, # running 孤儿任务巡检:容器重启/进程被杀后卡 running 的任务,20 分钟无更新则判失败 "cleanup-stale-running-tasks": { "task": "worker.cleanup_stale_running_tasks", "schedule": 300.0, # 每 5 分钟(秒) "options": {"expires": 240}, }, }