流式输出场景下最典型的问题就是连接莫名其妙地断了:前端正在接收SSE推送的数据,突然卡住;WebSocket聊天页面挂着不动了,过几分钟才发现早已掉线。这类问题的根源往往不在业务代码,而在网络链路、代理服务器的超时策略以及TCP连接本身的空闲回收机制。本文从掉线原因入手,重点讲心跳包机制的设计与实现,再给出几套断线重连策略的完整代码。

流式输出为什么会频繁中断
先搞清楚原因,才能对症下药。流式连接中断通常有以下几类。
第一类是中间代理的超时回收。Nginx默认对上游的proxy_read_timeout是60秒,如果60秒内没有任何字节从服务端发往客户端,代理就会主动切断连接。SSE这种长时间只收数据不回数据的场景特别容易中招。CDN、负载均衡器、云服务商的网关都有类似策略,超时时间从30秒到几分钟不等。
第二类是NAT设备的连接表过期。家用路由器、运营商网关会维护一张连接映射表,空闲超过一定时间(常见为几分钟)就把条目删掉。条目被删除后,客户端和服务端都不会收到任何通知,连接表面上看还在,实际上数据已经发不过去了,这就是典型的半开连接。
第三类是客户端网络切换。手机从WiFi切到4G、电梯里信号抖动,都会导致底层IP变化,原有TCP连接直接失效。
这三类问题有一个共同点:单靠业务数据流本身无法感知连接状态。所以解决方案的核心思路就是在业务数据之外,定期发送一些很小的保活数据,也就是心跳包,让链路上的每一环都知道这条连接还活着。
心跳包机制的设计与实现
心跳包的本质很简单:定时发送一个小数据帧,刷新中间设备的超时计时器,同时利用发送或接收结果探测连接是否健康。设计心跳时需要回答三个问题:谁来发、发多频、怎么判定死亡。
心跳间隔怎么定。间隔要小于链路上最短的超时时间。如果Nginx的read timeout是60秒,心跳间隔设30秒就比较稳妥。一般建议在15到30秒之间,移动端网络环境差的话可以适当缩短,但也不要太短,否则心跳本身的流量和电量消耗会成为新问题。
心跳由谁发。通常由客户端发更容易实现,因为客户端知道自己什么时候空闲。服务端也可以主动发,比如SSE场景下很多实现是服务端每隔15秒推一个注释行保活,这是SSE协议自带的便利特性。
下面是服务端SSE心跳的示例:
# Python + Flask 实现带心跳的 SSE 接口
import time
from flask import Flask, Response
app = Flask(__name__)
@app.route('/stream')
def stream():
def generate():
while True:
# 有业务数据时推送数据
data = get_next_data()
if data:
yield f"data: {data}\n\n"
else:
# 空闲时发送注释行作为心跳,防止代理超时断开
yield ": heartbeat\n\n"
time.sleep(1)
return Response(generate(), mimetype='text/event-stream')
注释行以冒号开头,SSE客户端会自动忽略它,不影响业务解析,但字节确实走了一遍网络链路,代理的计时器被刷新了。
WebSocket场景下一般用协议自带的ping和pong控制帧,浏览器端没有直接API发送ping,所以常见做法是客户端发一个自定义的文本消息作为心跳。下面是浏览器端的心跳检测实现:
// 浏览器端 WebSocket 心跳与死亡判定
class HeartbeatWS {
constructor(url) {
this.url = url;
this.pingInterval = 20000; // 每20秒发一次心跳
this.pongTimeout = 5000; // 心跳发出后5秒没响应则判定断开
this.init();
}
init() {
this.ws = new WebSocket(this.url);
this.ws.onopen = () => this.startHeartbeat();
this.ws.onmessage = (e) => {
if (e.data === 'pong') {
// 收到心跳响应,清除超时定时器
clearTimeout(this.pongTimer);
this.missed = 0;
} else {
this.handleData(e.data);
}
};
this.ws.onclose = () => this.stopHeartbeat();
}
startHeartbeat() {
this.timer = setInterval(() => {
this.ws.send('ping');
this.pongTimer = setTimeout(() => {
// 心跳超时,主动关闭触发重连
this.ws.close();
}, this.pongTimeout);
}, this.pingInterval);
}
stopHeartbeat() {
clearInterval(this.timer);
clearTimeout(this.pongTimer);
}
handleData(data) { /* 业务处理 */ }
}
这段代码体现了心跳死亡判定的标准做法:发送心跳后启动一个超时定时器,收到响应则清除,超时未清除就认为连接已死,主动关闭并进入重连流程。注意这里不能用onclose事件做唯一判断,因为半开连接下onclose可能长时间不触发。
断线重连策略的选择与实现
心跳负责发现连接死了,重连负责让连接复活。重连策略的好坏直接影响服务端压力和用户体验。
最基础的策略是立即重连,断开后马上尝试。这种做法在服务端短暂故障时恢复最快,但如果服务端已经宕机,成千上万个客户端同时立即重连,会形成重连风暴,把刚重启的服务再次压垮。
更推荐的是指数退避加抖动。每次重连失败后,等待时间按倍数增长,比如1秒、2秒、4秒、8秒,同时叠加一个随机抖动,避免所有客户端在同一时刻发起重连。另外要设置等待上限,比如最大30秒,超过后保持匀速重试:
function connectWithRetry(url, onMessage) {
let retry = 0;
const maxDelay = 30000;
function schedule() {
const base = Math.min(1000 * Math.pow(2, retry), maxDelay);
const jitter = Math.random() * base * 0.3; // 30%随机抖动
const delay = base + jitter;
setTimeout(doConnect, delay);
}
function doConnect() {
const ws = new WebSocket(url);
ws.onopen = () => { retry = 0; }; // 连接成功后重置计数
ws.onmessage = (e) => onMessage(e.data);
ws.onclose = () => {
retry++;
schedule();
};
}
doConnect();
}
第三个要考虑的问题是重连后的状态恢复。流式输出场景下,断线可能发生在一条数据只传了一半的时候,重连成功后怎么接着传?常见方案有两种。
方案一是客户端记录已接收的偏移量,重连请求时带上lastEventId或自定义的游标参数,服务端从该位置继续推送。SSE协议原生支持Last-Event-ID请求头,配合服务端的消息序号就能实现断点续传:
@app.route('/stream')
def stream_resume():
# 客户端重连时会自动带上 Last-Event-ID 请求头
last_id = int(request.headers.get('Last-Event-ID', 0))
def generate():
# 从上次中断的消息序号之后继续推送
for i in range(last_id + 1, get_total()):
yield f"id: {i}\ndata: {get_message(i)}\n\n"
return Response(generate(), mimetype='text/event-stream')
方案二是整体重发,服务端不保存状态,重连后重新生成完整输出。实现简单,但数据量大时会浪费带宽和算力,对生成式AI这类按token计费的流式接口来说成本明显。所以生产环境中,涉及长输出、高价值的流式接口,优先做状态恢复。
落地时的几个坑
实际部署中还有几个容易踩的坑值得提醒。一是别忘了调代理配置,如果Nginx的proxy_read_timeout比心跳间隔还短,心跳发得再勤也没用,需要同时调整类似proxy_read_timeout 120s这样的参数。二是浏览器页面切到后台时,定时器会被节流,心跳间隔可能被拉长到一分钟以上,可以考虑监听visibilitychange事件,页面回到前台时立即发一次心跳校验连接。三是重连逻辑要区分正常关闭和异常关闭,用户主动退出时不应触发自动重连,可以通过业务层的关闭标志位控制。
综合来看,一套稳定的流式输出方案就是心跳加指数退避重连加状态恢复三件套:心跳负责保活和探活,退避重连负责故障恢复,状态恢复负责数据不丢。三部分都不复杂,但缺了任何一块,长连接的体验都会大打折扣。