From f7f600d0910464dfccadc1f7fbb07b919d1e18f9 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Wed, 30 Sep 2026 15:30:36 +0800 Subject: [PATCH] fix(api): use celery_app.send_task() for viral_video enqueue instead of direct worker import The viral_video API route module had 4 places that did `from worker_app.tasks.viral_video import ` and then .delay(). In the API container the worker_app package is not installed, so calling generate/retry/confirm-intent/analyze-style would raise ModuleNotFoundError: No module named 'worker_app'. Switched all four call sites to celery_app.send_task() string-dispatch (matches the existing pattern in app/core/task_enqueue.py for worker.generate_video). Verified task names match @shared_task(name=...) declared in worker_app/tasks/viral_video.py: - worker.run_viral_video_pipeline - worker.resume_viral_video_pipeline - worker.run_video_style_analysis --- apps/api/app/api/routes/viral_video.py | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/apps/api/app/api/routes/viral_video.py b/apps/api/app/api/routes/viral_video.py index e9936203b..4444814df 100644 --- a/apps/api/app/api/routes/viral_video.py +++ b/apps/api/app/api/routes/viral_video.py @@ -34,6 +34,7 @@ from packages.adapters.sqlalchemy_impl.viral_video_repository import ( SQLAlchemyViralVideoJobRepository, SQLAlchemyViralVideoStyleTemplateRepository, ) +from app.core.celery_app import celery_app from packages.domain.viral_video import ViralVideoStatus logger = logging.getLogger(__name__) @@ -122,9 +123,7 @@ def create_viral_video( # 入队 Celery 任务 try: - from worker_app.tasks.viral_video import run_viral_video_pipeline - - run_viral_video_pipeline.delay(job.id) + celery_app.send_task("worker.run_viral_video_pipeline", args=[job.id]) logger.info("[爆款视频] 任务已入队: job_id=%s user_id=%s", job.id, job.user_id) except Exception as e: logger.error("[爆款视频] 入队失败: %s", e, exc_info=True) @@ -210,9 +209,7 @@ def retry_viral_video_job( # 重新入队 try: - from worker_app.tasks.viral_video import run_viral_video_pipeline - - run_viral_video_pipeline.delay(job.id) + celery_app.send_task("worker.run_viral_video_pipeline", args=[job.id]) logger.info("[爆款视频] 重试入队: job_id=%s retry_count=%d", job.id, job.retry_count) except Exception as e: logger.error("[爆款视频] 重试入队失败: %s", e, exc_info=True) @@ -249,9 +246,7 @@ def confirm_intent( # 从断点恢复 Celery 任务 try: - from worker_app.tasks.viral_video import resume_viral_video_pipeline - - resume_viral_video_pipeline.delay(job.id) + celery_app.send_task("worker.resume_viral_video_pipeline", args=[job.id]) logger.info("[爆款视频] 意图确认,恢复流水线: job_id=%s", job.id) except Exception as e: logger.error("[爆款视频] 恢复流水线失败: %s", e, exc_info=True) @@ -284,9 +279,7 @@ def analyze_style( # 入队风格分析任务 try: - from worker_app.tasks.viral_video import run_video_style_analysis - - run_video_style_analysis.delay(job.id) + celery_app.send_task("worker.run_video_style_analysis", args=[job.id]) logger.info("[爆款视频] 风格分析入队: job_id=%s", job.id) except Exception as e: logger.error("[爆款视频] 风格分析入队失败: %s", e, exc_info=True)