如何在FastApi中基于SQLAlchemy异步监听PostgreSQL Notify通道?
在FastAPI异步环境中用SQLAlchemy实现PostgreSQL Notify/Listen监听
核心思路
SQLAlchemy本身没有直接封装PostgreSQL的NOTIFY/LISTEN异步API,但可以通过它的异步连接获取底层的psycopg异步连接,结合FastAPI的后台任务和WebSocket能力实现实时聊天。下面是具体实现步骤:
1. 依赖准备
确保安装所需包:
pip install fastapi sqlalchemy[asyncio] psycopg[async] uvicorn python-multipart
2. 数据库连接配置
用SQLAlchemy异步引擎创建连接池和会话工厂:
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker # 替换为你的数据库URL DATABASE_URL = "postgresql+asyncpg://username:password@localhost/chat_db" # 创建异步引擎 engine = create_async_engine(DATABASE_URL, echo=True) # 异步会话工厂 AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
3. 实现PostgreSQL通道监听
创建独立的后台监听任务,复用SQLAlchemy的异步连接获取底层psycopg连接:
import asyncio from sqlalchemy.ext.asyncio import AsyncConnection from fastapi import WebSocket # WebSocket连接管理器,用于广播消息给前端 class WebSocketManager: def __init__(self): self.active_connections: set[WebSocket] = set() async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.add(websocket) def disconnect(self, websocket: WebSocket): self.active_connections.remove(websocket) async def broadcast(self, message: str): for connection in self.active_connections: await connection.send_text(message) manager = WebSocketManager() # 全局变量保存监听任务,用于关闭时清理 listen_task: asyncio.Task | None = None async def listen_to_postgres_channel(channel_name: str): # 获取独立的异步连接(不放入连接池,保持长连接监听) async with engine.connect() as conn: # 获取底层的psycopg异步连接 psycopg_conn = await conn.get_raw_connection() async with psycopg_conn.cursor() as cur: await cur.execute(f"LISTEN {channel_name};") print(f"开始监听PostgreSQL通道: {channel_name}") # 循环非阻塞接收通知 while True: # 1秒超时轮询,避免死锁 notification = await psycopg_conn.notifies.get(timeout=1) if notification: # 广播消息给所有WebSocket连接 await manager.broadcast(notification.payload)
4. 集成FastAPI生命周期和WebSocket端点
在FastAPI启动时启动监听任务,并提供WebSocket接口供前端连接:
from fastapi import FastAPI, WebSocket, WebSocketDisconnect app = FastAPI() # 启动时启动监听任务 @app.on_event("startup") async def startup_event(): global listen_task listen_task = asyncio.create_task(listen_to_postgres_channel("chat_messages")) # 关闭时清理监听任务 @app.on_event("shutdown") async def shutdown_event(): if listen_task: listen_task.cancel() try: await listen_task except asyncio.CancelledError: print("监听任务已取消") # WebSocket连接端点 @app.websocket("/ws/{client_id}") async def websocket_endpoint(websocket: WebSocket, client_id: str): await manager.connect(websocket) try: while True: # 接收前端发送的消息 data = await websocket.receive_text() # 将消息发送到PostgreSQL通道 async with AsyncSessionLocal() as session: # 注意转义单引号避免SQL注入 escaped_data = data.replace("'", "''") await session.execute(f"NOTIFY chat_messages, '{escaped_data}';") await session.commit() except WebSocketDisconnect: manager.disconnect(websocket)
关键注意事项
- 监听用的连接必须独立,不能与请求会话共用,否则会阻塞正常请求处理
- 发送
NOTIFY时必须转义payload中的单引号,避免SQL注入风险 - 若需要监听多个通道,可启动多个独立的监听任务
- 部署时建议使用单worker模式,多worker需额外处理进程间消息同步(小型应用单worker足够)
内容的提问来源于stack exchange,提问作者NothinPersonal
相关产品推荐
相关产品推荐

