使用psycopg3创建异步数据库连接的问题求助
Psycopg3 异步迁移问题解决指南
1. 「不能在函数外使用await」错误解决
错误代码:
con = await psycopg.AsyncConnection.connect(host="HOSTNAME", port="PORT", database="DATABASE", user="USER", password="PASSWORD") cursor = await con.cursor()
原因:await关键字仅能在async修饰的函数内部使用,全局作用域或普通函数中直接调用会触发错误。
解决方法:
将异步连接逻辑封装到async函数中,通过asyncio.run()启动执行:
import asyncio import psycopg async def main(): con = await psycopg.AsyncConnection.connect( host="HOSTNAME", port="PORT", database="DATABASE", user="USER", password="PASSWORD" ) cursor = await con.cursor() # 执行数据库操作... await con.close() asyncio.run(main())
2. 是否每个读写函数都要重复嵌套async with块?
不需要重复编写嵌套结构,可通过封装工具函数或使用连接池复用连接逻辑,减少冗余代码。
封装工具函数示例
编写通用执行函数,自动处理连接和游标生命周期:
async def execute_query(query, params=None): async with await psycopg.AsyncConnection.connect( host="HOSTNAME", port="PORT", database="DATABASE", user="USER", password="PASSWORD" ) as conn: async with conn.cursor() as cur: await cur.execute(query, params or ()) # 根据语句类型返回结果 if query.strip().upper().startswith("SELECT"): return await cur.fetchall() else: await conn.commit() return None
业务函数直接调用即可:
async def check_guild(guild_id): result = await execute_query( "SELECT guild_id, guild_name, su_id FROM guild WHERE guild_id = %s", [guild_id] ) return result[0] if result else None
使用连接池优化(高频场景推荐)
通过AsyncConnectionPool复用连接,避免重复创建连接的开销:
# 全局初始化连接池(程序启动时执行一次) pool = await psycopg.AsyncConnectionPool( host="HOSTNAME", port="PORT", database="DATABASE", user="USER", password="PASSWORD", min_size=2, max_size=10 ) async def execute_with_pool(query, params=None): async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute(query, params or ()) if query.strip().upper().startswith("SELECT"): return await cur.fetchall() else: await conn.commit() return None
3. AsyncConnectionPool 查询出现「AttributeError: 'AsyncConnection' object has no attribute 'fetchone'」
原因:fetchone是AsyncCursor游标对象的方法,错误地在AsyncConnection连接对象上调用了该方法。
解决方法:从连接中获取游标,通过游标执行查询并调用fetchone/fetchall:
async def check_guild_with_pool(guild_id): async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute( "SELECT guild_id, guild_name, su_id FROM guild WHERE guild_id = %s", [guild_id] ) return await cur.fetchone() # 游标对象调用fetchone
迁移后的完整业务函数示例(基于连接池)
对应原psycopg2同步代码,适配psycopg3异步版本:
import asyncio import psycopg import logging # 初始化连接池 async def init_pool(): global pool pool = await psycopg.AsyncConnectionPool( host="HOSTNAME", port="PORT", database="DATABASE", user="USER", password="PASSWORD", min_size=2, max_size=10 ) async def check_guild(guild_id): async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute( "SELECT guild_id, guild_name, su_id FROM guild WHERE guild_id = %s", [guild_id] ) return await cur.fetchone() async def config_raffle(guild_id, channel_id, channel_name, channel_cat_id, token, token_id, default_address, su_id, fee): try: async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute(""" INSERT INTO raffle_config (guild_id, channel_id, channel_name, channel_cat_id, token, default_token, default_address, su_id, fee) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (guild_id) DO UPDATE SET channel_id = EXCLUDED.channel_id, channel_name = EXCLUDED.channel_name, channel_cat_id = EXCLUDED.channel_cat_id, token = EXCLUDED.token, default_token = EXCLUDED.default_token, default_address = EXCLUDED.default_address, su_id = EXCLUDED.su_id, fee = EXCLUDED.fee """, (guild_id, channel_id, channel_name, channel_cat_id, token, token_id, default_address, su_id, fee)) await conn.commit() except Exception: logging.exception("Exception occurred") # async with块会自动回滚未提交事务,无需手动rollback print("Error: 25") return True # 程序入口 if __name__ == "__main__": asyncio.run(init_pool()) # 后续可调用业务函数
内容的提问来源于stack exchange,提问作者rolandj
相关产品推荐
相关产品推荐

