修复celery启动BUG稳定引入

This commit is contained in:
2026-06-15 15:13:08 +08:00
parent 1862a660f6
commit 2531276d15
2 changed files with 27 additions and 25 deletions
+7 -22
View File
@@ -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 任务。
"""
+20 -3
View File
@@ -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()