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

FastAPI永久运行Postgres监听后台任务避免随WebSocket断开取消

问题:FastAPI中PostgreSQL LISTEN后台任务随首个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连接绑定。

具体改法分三步:

  1. 修正ConnectionManager类的定义:类里直接写self.count_connections = 0是语法错误,实例属性要放在__init__方法里初始化,同时补全连接管理、广播的基础逻辑。
  2. 用FastAPI的lifespan生命周期(旧版可拆分使用startup/shutdown事件)管理监听任务:启动时创建独立的后台任务跑监听,不要把任务逻辑放在任何接口、WebSocket处理函数里;关闭时主动取消任务、释放数据库连接,避免资源泄漏。
  3. 给监听逻辑加异常重试和全局错误捕获:数据库连接断开时自动重连,不要让任务因为偶发错误直接退出,同时捕获任务内所有异常,避免出现"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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 10:12:27