You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

asyncpg/FastAPI项目数据库连接数超出连接池上限问题咨询

问题根因

你的代码存在两个核心错误,直接导致连接数超出预期:

1. 连接池初始化存在竞态条件,高并发下会重复创建多个连接池实例

当前connect方法没有加锁,当大量并发请求同时进入时,可能有数十个协程同时判断self._connection_pool is None,随后同时执行create_pool逻辑。每个连接池独立计数最大30个连接,只要同时创建4个连接池,就能直接打满你数据库100的连接上限。

2. 数据库操作方法逻辑分支错误,首次调用不会返回结果,易引发重复调用

你写的fetch_rows和execute方法,把获取连接执行查询的逻辑全放在了else分支里。也就是说第一次调用这两个方法时,因为连接池不存在,只会执行创建连接池的逻辑,不会执行任何查询操作也不会返回结果。如果你的测试脚本或者前端遇到无返回自动重试,会进一步加剧连接池重复创建的问题。

修复方案

修正分支逻辑,同时给连接池初始化加异步锁避免重复创建,修复后代码如下:

import asyncpg
import asyncio
class Database:
    def __init__(self, database: str, user: str, password: str, host: str, port: int = 5432):
        self.user = user
        self.password = password
        self.host = host
        self.port = port
        self.database = database
        self._connection_pool = None
        # 新增初始化锁,防止并发创建连接池
        self._init_lock = asyncio.Lock()

    async def connect(self):
        # 双重检查加锁,避免每次调用都抢锁影响性能
        if not self._connection_pool:
            async with self._init_lock:
                # 拿到锁之后二次检查,避免前一个协程已经完成连接池创建
                if not self._connection_pool:
                    try:
                        self._connection_pool = await asyncpg.create_pool(
                            min_size=1,
                            max_size=30,
                            command_timeout=300,
                            host=self.host,
                            port=self.port,
                            user=self.user,
                            password=self.password,
                            database=self.database,
                        )
                    except Exception as e:
                        print(e)
                        # 建议此处抛出异常不要静默失败,方便上层排查问题
                        raise
    
    async def disconnect(self):
        if self._connection_pool:
            try:
                await self._connection_pool.close()
                self._connection_pool = None
            except Exception as e:
                print(e)
                raise

    async def fetch_rows(self, query: str, *args):
        if not self._connection_pool:
            await self.connect()
        # 去掉else,创建完连接池后直接执行查询逻辑
        con = await self._connection_pool.acquire()
        try:
            result = await con.fetch(query, *args)
            return result
        except Exception as e:
            print(e)
            raise
        finally:
            await self._connection_pool.release(con)

    async def execute(self, query: str):
        if not self._connection_pool:
            await self.connect()
        # 去掉else,创建完连接池后直接执行查询逻辑
        con = await self._connection_pool.acquire()
        try:
            result = await con.execute(query)
            return result
        except Exception as e:
            print(e)
            raise
        finally:
            await self._connection_pool.release(con)

额外优化建议

  • 可以在FastAPI的启动事件中提前初始化连接池,不要等到请求进来才懒加载,从根源避免并发初始化问题
  • 可以给acquire方法加上超时时间,避免大量请求堆积等待连接时拖垮服务

内容的提问来源于stack exchange,提问作者Lostsoul

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 02:15:05