FastAPI永久运行Postgres监听后台任务避免随WebSocket断开取消
最小可复现示例
import asyncio import aiopg from fastapi import FastAPI, WebSocket dsn = "dbname=aiopg user=aiopg password=passwd host=127.0.0.1" app = FastAPI() class ConnectionManager: self.count_connections = 0 # 其他类函数和变量参考FastAPI官方文档 ... manager = ConnectionManager() async def send_and_receive_data(websocket: WebSocket): data = await websocket.receive_json() await websocket.send_text('Thanks for the message') # 后续处理接收到的数据 # 来自aiopg官方文档 # 该函数用于监听PostgreSQL通知 async def listen(conn): async with conn.cursor() as cur: await cur.execute("LISTEN channel") while True: msg = await conn.notifies.get() async def postgres_listen(): async with aiopg.connect(dsn) as listenConn: listener = listen(listenConn) await listener @app.get("/") def read_root(): return {"Hello": "World"} @app.websocket("/") async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) manager.count_connections += 1 if manager.count_connections == 1: await asyncio.gather( send_and_receive_data(websocket), postgres_listen() ) else: await send_and_receive_data(websocket)
问题描述
我正在使用Vue.js、FastAPI和PostgreSQL构建应用,本示例中尝试使用Postgres的listen/notify机制,结合WebSocket实现消息推送能力,应用除WebSocket端点外还包含大量常规HTTP端点。
需求是在FastAPI应用启动时运行一个永久存活的异步后台函数,该函数负责向所有WebSocket连接的客户端发送消息。即执行uvicorn main:app启动服务时,除了运行FastAPI应用本身,还需要同步启动后台函数postgres_listen(),当数据库表中新增数据行时,该函数会向所有在线WebSocket用户推送对应通知。
我了解可以使用asyncio.create_task()将任务放在on_*生命周期事件中启动,或是直接放在manager = ConnectionManager()实例化代码行之后启动,但该方案在我的场景下无法生效:只要发起任意HTTP请求(例如调用read_root()接口),就会抛出下文描述的相同错误。
当前采用的临时实现方案是:仅当首个客户端连接WebSocket时,才在websocket_endpoint()函数中启动postgres_listen()函数,后续任意客户端连接都不会重复运行/触发该函数。该方案原本运行正常,但当首个客户端/用户断开连接(例如关闭浏览器标签页)时,会立即抛出由psycopg2.OperationalError引发的GeneratorExit错误,错误日志如下:
Future exception was never retrieved future: <Future finished exception=OperationalError('Connection closed')> psycopg2.OperationalError: Connection closed Task was destroyed but it is pending! task: <Task pending name='Task-18' coro=<Queue.get() done, defined at /home/user/anaconda3/lib/python3.8/asyncio/queues.py:154> wait_for=<Future cancelled>>
该错误来自listen()函数,错误发生后asyncio的Task被取消,将无法再接收到来自数据库的通知。我已确认psycopg2、aiopg、asyncio本身不存在问题,核心疑问是应该将postgres_listen()函数放在什么位置,才能保证首个客户端断开连接后该监听任务不会被取消。我知道可以简单编写一个Python脚本作为首个客户端永久连接WebSocket,以此规避psycopg2.OperationalError异常,但该方案并不合理。
备注:
asyncio.shield()方案已测试,无法生效。
核心思路是把监听任务完全和WebSocket请求/连接生命周期解绑,只在FastAPI全局启动/关闭事件里管理这个后台任务,不要把它和任何单个请求、单个WebSocket连接绑定。
具体改法分三步:
- 修正
ConnectionManager类的定义:类里直接写self.count_connections = 0是语法错误,实例属性要放在__init__方法里初始化,同时补全连接管理、广播的基础逻辑。 - 用FastAPI的lifespan生命周期(旧版可拆分使用startup/shutdown事件)管理监听任务:启动时创建独立的后台任务跑监听,不要把任务逻辑放在任何接口、WebSocket处理函数里;关闭时主动取消任务、释放数据库连接,避免资源泄漏。
- 给监听逻辑加异常重试和全局错误捕获:数据库连接断开时自动重连,不要让任务因为偶发错误直接退出,同时捕获任务内所有异常,避免出现"Future exception was never retrieved"报错。
改完的可运行核心代码如下:
import asyncio import aiopg from contextlib import asynccontextmanager from fastapi import FastAPI, WebSocket, WebSocketDisconnect from typing import List dsn = "dbname=aiopg user=aiopg password=passwd host=127.0.0.1" class ConnectionManager: def __init__(self): self.count_connections = 0 self.active_connections: List[WebSocket] = [] async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.append(websocket) self.count_connections += 1 def disconnect(self, websocket: WebSocket): if websocket in self.active_connections: self.active_connections.remove(websocket) self.count_connections -= 1 async def broadcast(self, message: str): # 给所有在线连接发消息,自动剔除已经断开的无效连接 for connection in self.active_connections.copy(): try: await connection.send_text(message) except WebSocketDisconnect: self.disconnect(connection) manager = ConnectionManager() # 全局持有监听任务和数据库连接引用,统一做生命周期管理 listen_task = None db_conn = None async def listen_notify(): global db_conn # 内置重连逻辑,连接异常断开后自动恢复 while True: try: async with aiopg.connect(dsn) as conn: db_conn = conn async with conn.cursor() as cur: await cur.execute("LISTEN channel") while True: msg = await conn.notifies.get() # 收到数据库通知就广播给所有在线WebSocket客户端 await manager.broadcast(f"新通知: {msg.payload}") except Exception as e: print(f"数据库监听连接异常,5秒后重连: {str(e)}") await asyncio.sleep(5) finally: db_conn = None # 新版本FastAPI用lifespan统一管理启动/关闭逻辑,旧版本可拆分写@app.on_event("startup")和@app.on_event("shutdown") @asynccontextmanager async def lifespan(app: FastAPI): # 应用启动时创建独立后台任务,任务绑定应用主事件循环,和单个请求完全隔离 global listen_task listen_task = asyncio.create_task(listen_notify()) yield # 应用关闭时主动清理任务和连接 if listen_task and not listen_task.done(): listen_task.cancel() try: await listen_task except asyncio.CancelledError: pass if db_conn and not db_conn.closed: db_conn.close() app = FastAPI(lifespan=lifespan) async def send_and_receive_data(websocket: WebSocket): try: while True: data = await websocket.receive_json() await websocket.send_text('Thanks for the message') except WebSocketDisconnect: manager.disconnect(websocket) @app.get("/") def read_root(): return {"Hello": "World"} @app.websocket("/") async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) # WebSocket端点只处理当前连接的收发逻辑,完全不涉及监听任务的启停 await send_and_receive_data(websocket)
原有写法的问题根源
- 你把
postgres_listen()放在asyncio.gather()里和当前WebSocket的收发逻辑绑定,asyncio.gather()会等待所有传入协程完成,只要其中一个协程(比如客户端断开导致send_and_receive_data退出),gather就会把同组的其他所有协程全部取消,这就是首个客户端断开后监听任务直接崩溃的根本原因,asyncio.shield()无法阻挡gather上下文级的取消传播。 - 之前尝试在生命周期事件里启动任务报错,本质是没有捕获任务内部的异常,也没有正确处理aiopg连接的上下文生命周期,导致任务刚启动就因为未捕获异常直接退出。
- 上述方案里的监听任务是应用级全局任务,和任何单个HTTP请求、WebSocket连接都没有绑定关系,不管多少个客户端连接、断开,都不会影响后台监听任务的运行,内置的自动重连逻辑也能覆盖数据库临时重启、网络波动等异常场景。
内容的提问来源于stack exchange,提问作者artemonsh

