如何优化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
相关产品推荐
相关产品推荐

