From 2531276d1556ef89f455f8ed2fd2097e52f4c7d2 Mon Sep 17 00:00:00 2001 From: GinHa <15201596918@163.com> Date: Mon, 15 Jun 2026 15:13:08 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8Dcelery=E5=90=AF=E5=8A=A8BUG?= =?UTF-8?q?=E7=A8=B3=E5=AE=9A=E5=BC=95=E5=85=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- video-gen-api/app/tasks/__init__.py | 29 +++++++-------------------- video-gen-api/app/tasks/celery_app.py | 23 ++++++++++++++++++--- 2 files changed, 27 insertions(+), 25 deletions(-) diff --git a/video-gen-api/app/tasks/__init__.py b/video-gen-api/app/tasks/__init__.py index 81a62333..8f013790 100644 --- a/video-gen-api/app/tasks/__init__.py +++ b/video-gen-api/app/tasks/__init__.py @@ -1,24 +1,9 @@ -"""Celery task module imports. +"""Celery task package. -Celery autodiscover imports ``app.tasks``; importing task modules here ensures -custom named tasks are registered when workers start. -""" +任务模块注册统一由 app.tasks.celery_app.CELERY_TASK_IMPORTS 控制。 -import logging - -logger = logging.getLogger("video_gen") - -try: - from app.tasks import ( # noqa: F401 - generation_create_tasks, - generation_poll_tasks, - generation_download_tasks, - generation_recovery_tasks, - hot_opening_replicate_tasks, - shot_replicate_tasks, - shot_replicate_flow_tasks, - module_async_recovery_tasks, - ) -except Exception: - logger.exception("Celery 任务模块导入失败,worker 可能出现 unregistered task。") - raise +这里不要主动 import 子任务模块,避免以下问题: +1. Celery worker 启动时循环导入; +2. 新增任务文件后部分 worker 注册不完整; +3. app.tasks.__init__ 被普通业务代码 import 时意外加载全部 Celery 任务。 +""" \ No newline at end of file diff --git a/video-gen-api/app/tasks/celery_app.py b/video-gen-api/app/tasks/celery_app.py index e203f2c0..9ef946f7 100644 --- a/video-gen-api/app/tasks/celery_app.py +++ b/video-gen-api/app/tasks/celery_app.py @@ -10,6 +10,22 @@ from app.tasks.async_runner import close_loop, run_async logger = logging.getLogger("video_gen") +# 显式注册所有 Celery 任务模块,避免新增任务文件后 worker 启动时未注册任务。 +# 不再依赖 app.tasks.__init__ 内部 import,也不再依赖 autodiscover_tasks。 +CELERY_TASK_IMPORTS = ( + "app.tasks.generation_create_tasks", + "app.tasks.generation_poll_tasks", + "app.tasks.generation_download_tasks", + "app.tasks.generation_recovery_tasks", + "app.tasks.hot_opening_replicate_tasks", + "app.tasks.shot_replicate_tasks", + "app.tasks.shot_replicate_flow_tasks", + "app.tasks.module_async_recovery_tasks", + "app.tasks.user_oauth_tasks", + "app.tasks.cleanup", +) + + def _derive_redis_db(url: str, db_no: int) -> str: if not url: return url @@ -24,10 +40,11 @@ broker_url = settings.CELERY_BROKER_URL or (_derive_redis_db(settings.REDIS_URL, backend_url = settings.CELERY_RESULT_BACKEND or (_derive_redis_db(settings.REDIS_URL, 2) if settings.REDIS_URL else "") if broker_url: - celery_app = Celery("videogen") + celery_app = Celery("videogen", include=CELERY_TASK_IMPORTS) celery_app.conf.update( broker_url=broker_url, result_backend=backend_url or broker_url, + imports=CELERY_TASK_IMPORTS, task_serializer="json", accept_content=["json"], result_serializer="json", @@ -60,10 +77,10 @@ if broker_url: "generation.recover_download_tasks_once": {"queue": "gen_result_download"}, "generation.recover_generation_tasks_once": {"queue": "gen_result_download"}, "module_async.recover_module_async_tasks_once": {"queue": "gen_result_download"}, + "user_oauth.update_oauth_accounts": {"queue": "default"}, "app.tasks.cleanup.*": {"queue": "default"}, }, ) - celery_app.autodiscover_tasks(["app.tasks"]) else: celery_app = None @@ -164,4 +181,4 @@ def on_worker_process_shutdown(**kwargs): except Exception: pass finally: - close_loop() + close_loop() \ No newline at end of file