Python Flask中Server-Sent Events的并发问题及解决咨询
问题根源
你当前的实现里,message_list.pop(0)是消费型操作——消息被第一个客户端取走后就从列表中删除了,后续客户端自然拿不到这条消息。线程锁只是防止多个线程同时修改列表,但解决不了消息被移除的核心问题。
修复方案
下面提供两种可行的修复思路,根据你的业务需求选择:
方案一:基于消息ID的持久化推送(支持重连补收)
这种方案给每条消息分配唯一ID,每个客户端跟踪自己最后收到的ID,服务器只发送客户端未接收过的消息,不会删除消息列表中的内容,确保所有客户端都能拿到全量消息。
后端代码修改
from random import random from threading import Lock import time import json from flask import Flask, Response, request app = Flask(__name__) # 存储消息:[(消息ID, 内容), ...] messages = [] current_message_id = 0 message_lock = Lock() def event_stream(last_received_id): global current_message_id while True: new_messages = [] with message_lock: # 筛选出客户端未接收的消息 new_messages = [msg for msg in messages if msg[0] > last_received_id] # 更新客户端的最后接收ID(如果有新消息) if new_messages: last_received_id = new_messages[-1][0] # 逐个发送新消息 for msg_id, content in new_messages: # 用JSON格式携带消息ID和内容 yield f"data: {json.dumps({'id': msg_id, 'content': content})}\n\n" time.sleep(1) @app.route('/stream') def stream(): # 从查询参数获取客户端最后接收的ID,默认0(从未接收过) last_id = int(request.args.get('last_id', 0)) return Response(event_stream(last_id), mimetype="text/event-stream") @app.route("/new_message") def new_message(): global current_message_id with message_lock: current_message_id += 1 messages.append( (current_message_id, random()) ) # 限制消息列表长度,避免内存溢出(可选) if len(messages) > 100: messages.pop(0) return Response(status=204)
前端代码修改
需要跟踪最后收到的消息ID,重连时传递给服务器,确保不会遗漏消息:
let lastReceivedId = 0; let eventSource; function connectStream() { // 连接时带上最后接收的ID eventSource = new EventSource(`/stream?last_id=${lastReceivedId}`); eventSource.onmessage = function(event) { const msg = JSON.parse(event.data); console.log("收到消息:", msg.content); // 更新最后接收的ID lastReceivedId = msg.id; }; eventSource.onerror = function() { eventSource.close(); // 断开后3秒重连 setTimeout(connectStream, 3000); }; } // 初始化连接 connectStream();
方案二:发布订阅模式(仅推送在线客户端)
这种方案维护一个在线客户端的订阅列表,当有新消息时,直接推送给所有在线客户端,不存储历史消息,适合只需要实时通知在线用户的场景。
后端代码修改
from random import random from threading import Lock, Thread import time from flask import Flask, Response app = Flask(__name__) # 存储所有在线客户端的生成器 subscribers = [] subscribers_lock = Lock() def notify_all_subscribers(message): with subscribers_lock: # 遍历订阅者副本,避免迭代时列表被修改 for subscriber in list(subscribers): try: # 向客户端发送消息 subscriber.send(f"data: {message}\n\n") except Exception: # 客户端断开连接,移除订阅者 subscribers.remove(subscriber) def event_stream(): # 创建一个生成器用于向单个客户端发送消息 def client_generator(): yield "" # 发送空数据确保连接建立 gen = client_generator() next(gen) # 预激生成器 # 将生成器加入订阅列表 with subscribers_lock: subscribers.append(gen) try: while True: # 挂起生成器,等待新消息 yield next(gen) finally: # 客户端断开,从订阅列表移除 with subscribers_lock: if gen in subscribers: subscribers.remove(gen) @app.route('/stream') def stream(): return Response(event_stream(), mimetype="text/event-stream") @app.route("/new_message") def new_message(): msg = random() # 启动线程推送消息,避免阻塞请求 Thread(target=notify_all_subscribers, args=(msg,)).start() return Response(status=204)
前端代码
无需修改,保持你原来的代码即可:
var source = new EventSource("/stream"); source.onmessage = function(event) { console.log(event.data); };
方案对比
| 方案 | 优点 | 缺点 |
|---|---|---|
| 基于消息ID的持久化推送 | 支持客户端重连补收错过的消息,逻辑清晰 | 需要维护消息列表,需注意内存占用 |
| 发布订阅模式 | 实时性高,无需存储历史消息 | 客户端断开后无法获取期间的消息,需处理订阅者清理 |
内容的提问来源于stack exchange,提问作者Charles Dupont
相关产品推荐
相关产品推荐

