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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 09:31:43