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 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}, }, }