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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 07:40:56