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

