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

Flask结合PubSub实现SSE时多订阅者被阻塞问题排查

问题分析与解决方案

我用Flask开发实时数据表页面,数据来自PostgreSQL数据库,通过PubSub传递新增数据,再用SSE推送到前端。单个窗口运行正常,但多窗口模拟多用户时,只有最新加载的窗口能接收SSE,试过Flask线程模式、gevent、gunicorn等方案都没解决。

现有代码

main.py

from gevent import monkey
monkey.patch_all()
import psycopg2, pgpubsub
from flask import Flask, render_template, Response
from gevent.pywsgi import WSGIServer
import time


app = Flask(__name__)

pubsub = pgpubsub.connect(database="postgres"....)


def get_connection():
    try:
        return psycopg2.connect(database="postgres.....")
    except Exception as e:
        return f"Error connecting to DB: {e}"


@app.route('/')
def home():
    return render_template('index.html')

@app.route('/events')
def events():
    def update_pusher():
        print('Started listening')
        pubsub.listen('data_changed')
        while True:
            for event in pubsub.events(yield_timeouts=True):
                if event is None:
                    pass
                else:
                    yield f"data: {event.payload}\nevent: online\n\n"
                    time.sleep(0.01)
    return Response(
        response=update_pusher(),
        mimetype='text/event-stream'
    )

if __name__ == "__main__":
  http_server = WSGIServer(("localhost", 5003), app)
  http_server.serve_forever()

index.html

<p>This list is populated by server side events.</p>
<ul id="list"></ul>
<script>
    var eventSource = new EventSource("/events")
    eventSource.addEventListener("online", function(e) {
      // console.log(e.data.color)
      data = JSON.parse(e.data)
      const li = document.createElement('li')
      li.innerText = data
      list.appendChild(li)
    }, false)
</script>

问题根源

当前的pubsub是全局单例对象,所有请求共用同一个监听实例。PostgreSQL的LISTEN/NOTIFY机制是基于单个连接的,当新请求进来调用pubsub.listen('data_changed')时,会覆盖旧连接的监听上下文,导致之前的SSE连接无法再接收事件。

修复方案

每个SSE连接都需要创建独立的PubSub实例,而非共用全局对象。修改events视图函数,在生成器内部初始化专属的PubSub连接:

修改后的main.py关键代码

@app.route('/events')
def events():
    def update_pusher():
        print('Started listening')
        # 每个SSE连接创建独立的PubSub实例
        local_pubsub = pgpubsub.connect(database="postgres"....)
        local_pubsub.listen('data_changed')
        try:
            while True:
                for event in local_pubsub.events(yield_timeouts=True):
                    if event is not None:
                        yield f"data: {event.payload}\nevent: online\n\n"
                    # 用gevent.sleep替代time.sleep,避免阻塞协程
                    from gevent import sleep
                    sleep(0.01)
        finally:
            # 连接关闭时清理资源
            local_pubsub.unlisten('data_changed')
            local_pubsub.close()
    return Response(
        response=update_pusher(),
        mimetype='text/event-stream',
        headers={
            'Cache-Control': 'no-cache',
            'Connection': 'keep-alive',
            'X-Accel-Buffering': 'no'  # 禁用反向代理缓冲,确保实时推送
        }
    )

额外优化点

  • 确保PostgreSQL用户拥有LISTEN/NOTIFY权限,默认权限通常满足需求,但如果遇到权限错误需手动授权。
  • 若使用反向代理(如Nginx),需配置禁用缓冲,否则SSE推送会被延迟。
  • 可添加连接超时处理,避免无效连接长期占用资源。

内容的提问来源于stack exchange,提问作者kirstendo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 00:35:19