8210806632
1. 修复 Worker 启动入口: celery_app → worker_app.celery_app - 旧 celery_app.py 不导入任何任务模块,导致 Worker 注册 0 个任务 - 删除旧版 apps/worker/celery_app.py,统一使用 worker_app/celery_app.py 2. 为缺少装饰器的任务补充 @celery_app.task: - classification.py: classify_asset() 添加装饰器 - generation.py: generate_video() 添加装饰器,修正签名匹配 API 调用方式 3. 确保 voice_extraction 任务被正确注册: - 添加 voice_extraction 到 celery_app imports - 修复 voice_extraction.py 中错误的相对导入 (.celery_app → worker_app.celery_app) - 修复 dedup.py 中指向已删除模块的导入 4. 修复 worker_app/celery_app.py: - 添加 broker_connection_retry_on_startup=True - imports 中添加 voice_extraction 和 dedup 模块 5. 修复 Dockerfile: - CMD 改为 celery -A worker_app.celery_app - 添加非 root 用户 celery 运行 Worker 6. 新建 packages/shared/config.py 和 storage.py 兼容层 - 为 worker 任务模块提供统一的 config/storage 访问入口
17 lines
556 B
Python
17 lines
556 B
Python
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",
|
|
"apps.worker.video_processing.dedup",
|
|
)
|