From eb0a4005e1254d41c76ca688458c181e0ccaf368 Mon Sep 17 00:00:00 2001 From: chenjw28 <792430652@qq.com> Date: Sat, 26 Sep 2026 23:43:54 +0800 Subject: [PATCH] =?UTF-8?q?fix(db):=20celery=20=E8=A7=92=E8=89=B2=E6=94=B9?= =?UTF-8?q?=E7=94=A8=20NullPool=E2=80=94=E2=80=94=E8=BF=9E=E6=8E=A5?= =?UTF-8?q?=E8=B7=A8=E4=BA=8B=E4=BB=B6=E5=BE=AA=E7=8E=AF=E5=A4=8D=E7=94=A8?= =?UTF-8?q?=E8=87=B4=E4=BB=BB=E5=8A=A1=E5=85=A8=E6=8C=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit celery_tasks 经 asyncio.run 每任务新建事件循环,而 asyncpg 连接绑定 创建它的循环。db_manager 引擎只连一次(is_connected 短路),QueuePool 把上一循环的连接缓存着流入新循环:任务报 "got Future attached to a different loop" 全部 failed,连接销毁时再报 "Event loop is closed"(ForkPoolWorker 日志连续复现)。 celery_tasks 对 redis 已有同因处理(每任务 reconnect()),DB 侧漏了。 改法:connect(role="celery") 用 NullPool——连接不缓存,每次 checkout 在当前循环新建、用完即关。web 入口(单循环 uvicorn)QueuePool 不变。 已实测:同进程三个连续 asyncio.run 循环共享引擎,全部查询成功。 Co-Authored-By: Claude Code --- src/shared/database/database.py | 46 +++++++++++++++++++++++---------- 1 file changed, 32 insertions(+), 14 deletions(-) diff --git a/src/shared/database/database.py b/src/shared/database/database.py index 77028a9..8f7a959 100644 --- a/src/shared/database/database.py +++ b/src/shared/database/database.py @@ -53,16 +53,31 @@ class DatabaseManager: return try: - pool_cfg = _get_pool_config(role) - # 创建异步引擎 - self.engine = create_async_engine( - database_url, - echo=settings.DEBUG, - pool_size=pool_cfg["pool_size"], - max_overflow=pool_cfg["max_overflow"], - pool_recycle=3600, - pool_pre_ping=True, # 自动检测失效连接,避免 PG 断连报错 - ) + if role == "celery": + # Celery 场景必须用 NullPool:celery_tasks 经 asyncio.run 每任务 + # 新建事件循环,而 asyncpg 连接绑定创建它的循环——QueuePool 会把 + # 上一循环的连接缓存着流入新循环,报 + # "got Future attached to a different loop"(任务全挂)且销毁时 + # "Event loop is closed"。NullPool 不缓存:每次 checkout 在当前 + # 循环新建连接、用完即关。局域网建连 ~1ms,对分钟级分析任务可忽略。 + # (与 redis_task_manager 每任务 reconnect() 是同一循环绑定问题) + from sqlalchemy.pool import NullPool + self.engine = create_async_engine( + database_url, + echo=settings.DEBUG, + poolclass=NullPool, + ) + else: + pool_cfg = _get_pool_config(role) + # 创建异步引擎 + self.engine = create_async_engine( + database_url, + echo=settings.DEBUG, + pool_size=pool_cfg["pool_size"], + max_overflow=pool_cfg["max_overflow"], + pool_recycle=3600, + pool_pre_ping=True, # 自动检测失效连接,避免 PG 断连报错 + ) # 创建异步会话工厂 self.async_session = async_sessionmaker( @@ -76,10 +91,13 @@ class DatabaseManager: await conn.execute(text("SELECT 1")) self.is_connected = True - logger.info( - "数据库连接成功 (pool_size=%d, max_overflow=%d)", - pool_cfg["pool_size"], pool_cfg["max_overflow"], - ) + if role == "celery": + logger.info("数据库连接成功 (role=celery, pool=NullPool)") + else: + logger.info( + "数据库连接成功 (role=%s, pool_size=%d, max_overflow=%d)", + role, pool_cfg["pool_size"], pool_cfg["max_overflow"], + ) except Exception as e: logger.error(f"数据库连接失败: {e}")