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}")