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

如何基于FastAPI WebSocket实现Postgres LISTEN异步非阻塞功能?

解决FastAPI WebSocket结合Postgres LISTEN的异步阻塞问题

嘿,我看你遇到的问题挺典型——用同步的psycopg2在异步FastAPI路由里折腾,难怪会踩坑。咱们一步步理清楚:

为什么会报错?

你看到的object psycopg2.extensions.Notify can't be used in the 'await' expression,直接原因是conn.notifies.pop(0)返回的是普通对象,不是异步可等待(awaitable)对象,所以不能加await前缀。

但更深层的问题是:psycopg2是完全同步的库,在FastAPI的异步上下文里用它的同步方法(比如conn.poll())会阻塞整个事件循环,这就失去了异步Web应用的优势,甚至会导致其他请求被卡住。

正确解决方案:用异步版Psycopg

psycopg团队早就推出了异步版本的库——psycopg[async](也就是Psycopg 3的异步API),专门适配异步场景。咱们用它重写你的代码:

第一步:安装异步依赖

先把异步版psycopg装上:

pip install psycopg[async]

第二步:重写WebSocket路由代码

from fastapi import APIRouter, WebSocket, WebSocketDisconnect
import psycopg
from psycopg import sql

router = APIRouter()

@router.websocket("/pg_notify")
async def get_notifications(websocket: WebSocket):
    await websocket.accept()
    # 用异步方式连接Postgres
    async with await psycopg.AsyncConnection.connect(
        "your_postgres_connection_string_here",
        autocommit=True  # LISTEN必须在自动提交模式下执行
    ) as conn:
        # 创建异步游标并执行LISTEN命令
        async with conn.cursor() as cur:
            await cur.execute(sql.SQL("LISTEN {};").format(sql.Identifier("channel")))
        
        try:
            while True:
                # 异步等待Postgres通知,超时30秒避免无限阻塞
                notify = await conn.wait_for_notify(timeout=30)
                if notify:
                    # 收到通知后推送给WebSocket客户端
                    await websocket.send_text(f"收到更新:{notify.payload}")
                else:
                    # 超时发送心跳,保持连接活跃
                    await websocket.send_text("心跳:等待更新中...")
        except WebSocketDisconnect:
            print("客户端主动断开WebSocket连接")
        except Exception as e:
            print(f"发生错误:{str(e)}")
            await websocket.close()
        finally:
            # 主动取消监听(可选,连接关闭时Postgres会自动清理)
            async with conn.cursor() as cur:
                await cur.execute(sql.SQL("UNLISTEN {};").format(sql.Identifier("channel")))

关键代码说明

  • psycopg.AsyncConnection.connect:创建异步数据库连接,必须用await等待连接建立
  • autocommit=True:LISTEN命令需要在非事务上下文执行,所以必须开启自动提交
  • await conn.wait_for_notify(timeout=...):异步等待通知,不会阻塞FastAPI的事件循环,超时后返回None,方便处理心跳逻辑
  • sql.SQL + sql.Identifier:安全拼接SQL标识符(比如channel名称),避免SQL注入风险
  • 异常处理:专门捕获WebSocketDisconnect处理客户端主动断开的情况,确保资源正确清理

额外注意事项

  • 确保你的Postgres版本在9.0以上(基本都满足),支持LISTEN/NOTIFY机制
  • 如果需要监听多个channel,在同一个连接上执行多个LISTEN命令即可,wait_for_notify会接收所有channel的通知
  • 生产环境建议用AsyncConnectionPool管理连接池,避免频繁创建/销毁连接的开销

内容的提问来源于stack exchange,提问作者wkamer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 12:02:32