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

Asyncpg.pool创建连接数超出max_size的异常问题咨询

环境信息

  • asyncpg版本: 0.29.0
  • PostgreSQL版本: "PostgreSQL 16.0 (Debian 16.0-1.pgdg120+1) on aarch64-unknown-linux-gnu, compiled by gcc (Debian 12.2.0-14) 12.2.0, 64-bit"
  • 是否使用PostgreSQL SaaS?若使用,是哪家?能否在本地PostgreSQL复现问题?: 否,使用Docker容器部署
  • Python版本: 3.11.5
  • 平台: MacBook-Air Darwin Kernel Version 21.1.0: Wed Oct 13 17:33:24 PDT 2021; root:xnu-8019.41.5~1/RELEASE_ARM64_T8101 arm64
  • 是否使用pgbouncer?: 否
  • 是否通过pip安装asyncpg?: 是
  • 若本地编译asyncpg,使用的Cython版本是多少?: 不适用
  • 能否在asyncio和uvloop下均复现问题?: 否

问题描述

实现了Database单例用于写入数据,使用asyncio.gather()并发执行1000次写入时,数据库显示的活跃连接数远超asyncpg.pool的max_size(测试时数据库有857个活跃连接,但连接池仅显示62个活跃连接)。使用uvloop执行相同操作时,任务数超过池大小会直接崩溃,抛出ConnectionResetError: [Errno 54] Connection reset by peer错误。请问这属于连接池的正常行为吗?

代码示例

数据库代码

import os
import asyncpg
import asyncio

class Database:
    _instance = None
    _pool = None
    _lock = asyncio.Lock()
    db_params = { 
                'host': os.getenv('DATABASE_HOST'),
                'port': os.getenv('DATABASE_PORT'),
                'database': os.getenv('DATABASE_NAME'),
                'user': os.getenv('DATABASE_USER'),
                'password': os.getenv('DATABASE_PASSWORD')
            }

    def __new__(cls, *args, **kwargs):
        if cls._instance is None:
            cls._instance = super(Database, cls).__new__(cls)
        return cls._instance

    @classmethod
    async def get_pool(cls):
        if cls._pool is None:
            async with cls._lock:
                if cls._pool is None:
                    cls._pool = await asyncpg.create_pool(**cls.db_params, min_size=1, max_size=150)
        return cls._pool

    @classmethod
    async def write(cls, result):
            pool = await cls.get_pool()
            try:
                    async with pool.acquire() as connection:
                        result = await connection.execute('''
                            INSERT INTO tables.results(
                                result
                            ) VALUES($1)
                        ''', result)
                        return
            except Exception as e:
                raise e

演示写入代码

async def fake_result(i):
    print(f'generating fake result {i}')
    await Database.write(i)
    return

async def run_functions_concurrently():
   tasks = [fake_result(i) for i in range(1000)]
   await asyncio.gather(*tasks)

def main():
    asyncio.run(run_functions_concurrently())

if __name__ == "__main__":
    main()

问题分析与解决方案

这不属于连接池的正常行为,问题出在你的单例连接池实现存在协程安全漏洞:

根本原因

当大量协程并发调用get_pool()时,由于没有同步机制,多个协程会同时检测到cls._pool为None,并各自执行await asyncpg.create_pool(),最终创建了多个独立的连接池。每个连接池都遵循自己的max_size=150限制,但多个池的连接数累加后,就会导致数据库的总活跃连接数远超单个池的限制。

uvloop下的ConnectionResetError也是由此引发:短时间内创建大量连接,超出了PostgreSQL或操作系统的连接处理能力,导致服务器主动重置连接。

修复方案

给get_pool()添加异步锁,结合双重检查锁机制,确保同一时间只有一个协程能创建连接池,避免重复创建:

  1. 在Database类中添加异步锁:_lock = asyncio.Lock()
  2. 修改get_pool()方法,加入锁保护和双重检查(如上述代码示例所示)

额外建议

  • 修复后可通过pool.get_stats()查看连接池状态,确认仅存在一个连接池,且活跃连接数不超过max_size。
  • 检查PostgreSQL的max_connections配置(默认通常为100),确保其值不小于连接池的max_size,避免数据库层面拒绝新连接。

内容的提问来源于stack exchange,提问作者no-to-mediocrity

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 06:06:17