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

如何优化asyncpg多查询异步执行效率?解决异步慢于同步问题

问题:异步PostgreSQL查询性能劣于同步,复用连接触发报错

现有实现代码

数据库操作类与连接池单例

from typing import Optional, Union, List, Dict, Any
import asyncpg

class PgDb:
    # noinspection SpellCheckingInspection
    def __init__(self, conn: asyncpg.connection.Connection):
        self.conn = conn
    async def select(self, sql: str, args: Union[list, Dict[str, Any]] = []) -> List[Dict[str, Any]]:
        sql, _args = self.__convert_placeholders(sql, args)
        return [dict(row) for row in await self.conn.fetch(sql, *_args)]

class DbPoolSingleton:
    db_pool: Optional[asyncpg.pool.Pool] = None

    @staticmethod
    async def create_pool():
        config = get_postgres_config()
        pool: asyncpg.Pool = await asyncpg.create_pool(
            ...,
            min_size=30,
            max_size=40
        )
        print("Pool created")
        return pool

    @staticmethod
    async def get_pool() -> asyncpg.pool.Pool:
        if not DbPoolSingleton.db_pool:
            DbPoolSingleton.db_pool = await DbPoolSingleton.create_pool()
        return DbPoolSingleton.db_pool

    @staticmethod
    async def terminate_pool():
        (await DbPoolSingleton.get_pool()).terminate()
        DbPoolSingleton.db_pool = None
        print("Pool terminated")

测试代码与结果

同步测试

import asyncio
import datetime
from helpers.pg_rdb_helper import DbPoolSingleton, PgDb

async def test_synchronous():
    conn = await (await DbPoolSingleton.get_pool()).acquire()
    db = PgDb(conn)
    sql = """samplesql"""
    total_start = start = datetime.datetime.now()
    
    for i in range(20):
        start = datetime.datetime.now()
        rows = await db.select(sql)
        end  = datetime.datetime.now()
        print(f"{i}st query took: ", (end-start).total_seconds())
    total_end = datetime.datetime.now()
    print(f"total query took: ", (total_end-total_start).total_seconds())

执行结果:total query took: 2.131297

初始异步测试

async def test_asynchronous():
    db_pool = await DbPoolSingleton.get_pool()
    sql = """samplesql"""
    total_start = datetime.datetime.now()
    tasks = []
    for i in range(20):
        db = PgDb(await db_pool.acquire())
        task = asyncio.create_task(db.select(sql))
        tasks.append(task)
    await asyncio.gather(*tasks)
    total_end = datetime.datetime.now()
    print(f"total query took: ", (total_end-total_start).total_seconds())

执行结果:total query took: 2.721282

复用单连接的异步测试(报错)

async def test_asynchronous():
    db_pool = await DbPoolSingleton.get_pool()
    sql = """samplesql"""
    total_start = datetime.datetime.now()
    tasks = []
    async with db_pool.acquire() as conn:
        db = PgDb(conn)
        for i in range(20):
            task = asyncio.create_task(db.select(sql))
            tasks.append(task)
        await asyncio.gather(*tasks)
    total_end = datetime.datetime.now()
    print(f"total query took: ", (total_end-total_start).total_seconds())

报错信息:

asyncpg.exceptions._base.InterfaceError: cannot perform operation: another operation is in progress

问题原因

  • 初始异步版本耗时更高:每个任务单独调用db_pool.acquire()获取连接,频繁的连接获取/释放带来额外开销
  • 复用单连接报错:asyncpg的单个数据库连接是单操作串行的,同一连接上不能同时执行多个异步查询,否则会触发操作冲突

解决方案

优化后的异步测试代码

通过上下文管理器自动管理连接的获取与释放,避免手动操作的开销,同时确保每个任务使用独立的连接实现并发:

async def test_asynchronous_optimized():
    db_pool = await DbPoolSingleton.get_pool()
    sql = """samplesql"""
    total_start = datetime.datetime.now()
    
    # 封装单个查询的逻辑,自动处理连接的获取和释放
    async def execute_query():
        async with db_pool.acquire() as conn:
            db = PgDb(conn)
            return await db.select(sql)
    
    # 创建并发任务列表
    tasks = [asyncio.create_task(execute_query()) for _ in range(20)]
    await asyncio.gather(*tasks)
    
    total_end = datetime.datetime.now()
    print(f"total query took: ", (total_end-total_start).total_seconds())

额外优化建议

  • 连接池参数调优:当前连接池min_size=30,已覆盖20个并发任务的需求,无需扩容;若后续并发量提升,可适当调整max_size避免连接等待
  • 预编译SQL:对于重复执行的SQL,可提前预编译为asyncpg.PreparedStatement,减少数据库端的SQL解析开销
  • 优化占位符转换:确保PgDb.__convert_placeholders方法逻辑高效,避免不必要的字符串操作损耗性能

内容的提问来源于stack exchange,提问作者DFX Nguyễn

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 00:24:57