From 3f18b28227944494de97832fac24f362b767c42d Mon Sep 17 00:00:00 2001 From: wwwwwwwww <526125649@qq.com> Date: Fri, 14 Aug 2026 14:46:50 +0800 Subject: [PATCH] 1 --- video-gen-api/app/tasks/celery_app.py | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/video-gen-api/app/tasks/celery_app.py b/video-gen-api/app/tasks/celery_app.py index 1b8efc87..0aa3eaaf 100644 --- a/video-gen-api/app/tasks/celery_app.py +++ b/video-gen-api/app/tasks/celery_app.py @@ -1,3 +1,4 @@ +import asyncio import logging from celery import Celery @@ -331,6 +332,9 @@ def _setup_dynamic_beat_tasks(sender, **kwargs): 通过 @celery_app.on_after_configure.connect 在 Celery 配置完成后执行, 适用于 Worker 和 Beat 启动场景。 + + 注意:此信号运行在 Beat/Worker 主线程,不能使用 run_async(single_loop 在另一线程), + 这里用独立 event loop 同步运行 async 查询,避免 asyncpg Future 跨 loop 错误。 """ if celery_app is None: return @@ -345,16 +349,21 @@ def _setup_dynamic_beat_tasks(sender, **kwargs): result = await db.execute( select(ScheduledTask).where(ScheduledTask.is_active.is_(True)) ) - tasks = result.scalars().all() - return tasks + return result.scalars().all() + active_tasks = None try: - active_tasks = run_async(_load()) + # 在当前线程创建独立 loop 同步运行,不经过 single_loop + loop = asyncio.new_event_loop() + try: + active_tasks = loop.run_until_complete(_load()) + finally: + loop.close() except Exception: logger.exception("加载定时任务失败,跳过动态 Beat 注册") return - for task in active_tasks: + for task in active_tasks or []: schedule_val = _parse_schedule_to_celery(task.schedule) if schedule_val is None: logger.warning("定时任务 %s schedule 无效,跳过注册: %s", task.id, task.schedule)