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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 18:51:21