fix(db): celery 角色改用 NullPool——连接跨事件循环复用致任务全挂

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 <noreply@anthropic.com>
This commit is contained in:
2026-09-26 23:43:54 +08:00
parent 73de137779
commit eb0a4005e1
+32 -14
View File
@@ -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}")