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
相关产品推荐
相关产品推荐

