AI智能体的流式输出通常依赖SSE(Server-Sent Events)协议,将生成内容以数据块形式持续推送给前端。前端通过EventSource对象订阅事件流,每当服务端有新数据时,浏览器触发onmessage回调,UI实时更新。这种模式相比一次性HTTP请求,能大幅降低首字延迟,提升用户对智能体“正在思考并输出”的感知。然而,SSE连接本质上是长连接的HTTP响应,一旦中途断开,EventSource不会自动重新建立连接,导致流式输出彻底停止。更棘手的是,很多开发者错误地认为SSE自带重连功能,实际上浏览器只会触发onerror事件,并不会自动发起新的连接。针对Agent类应用,输出内容往往与任务状态强相关,中断后如果不做恢复,用户只能手动刷新,任务可能已经执行了一半却无法继续。因此,设计一套可靠的SSE重连机制,是保障AI智能体稳定性的关键环节。

问题定位:SSE连接中断的常见原因
要解决重连问题,首先得弄清楚连接为何中断。SSE运行在HTTP/1.1或HTTP/2之上,客户端发起一个普通的GET请求,服务端保持响应不结束,并不断写入数据。任何导致TCP连接关闭的因素都会让SSE终止。其中最常见的是网络层问题,比如用户切换Wi-Fi、移动网络波动、中间代理或负载均衡器的空闲超时。许多云服务商的负载均衡默认在60秒左右断开空闲连接,如果SSE没有心跳数据,就会被误判为死连接而切断。其次是服务端主动关闭,例如后端进程重启、部署新版本、内存溢出导致进程崩溃,或者代码中主动调用了response.end()。还有一类是客户端行为,例如用户关闭页面、浏览器标签页被系统休眠(尤其在移动端),或者前端代码主动调用了EventSource实例的close方法。除此之外,代理服务器缓冲和防火墙策略也可能干扰长连接,例如某些企业防火墙会阻断长时间无响应的HTTP流。
定位问题需要观察客户端表现和服务端日志。客户端在连接断开时会触发onerror事件,但该事件不提供具体的错误原因。开发者可以通过EventSource实例的readyState属性判断当前状态:0表示连接中,1表示连接已打开,2表示连接已关闭。服务端日志则能更清晰地看到连接何时被断开、由哪一方发起。为了更细致地监控,可以在服务端对SSE响应注册close事件,记录断开时的客户端IP、用户标识和已发送的事件数量。一个常见的坑是服务端在写入数据后返回了错误的Content-Type,或者没有设置Transfer-Encoding: chunked,导致代理层无法正确处理流式响应而提前结束。因此,排查阶段建议先确认响应头中Content-Type是否为text/event-stream,并且禁用任何缓冲中间件。
下面是一个最基础的SSE客户端实现,它只能接收消息,没有任何断线处理,可以作为问题复现的起点:
// 最简单的EventSource用法,不具备重连能力
const source = new EventSource('/api/agent/stream?task_id=123');
source.onmessage = function(event) {
const data = JSON.parse(event.data);
console.log('收到数据:', data);
// 将数据追加到UI,例如聊天消息体
appendMessage(data);
};
source.onerror = function(event) {
// 连接出错时触发,但不会自动重连
console.error('SSE连接出现错误,当前状态:', source.readyState);
// 这里需要自己实现重连逻辑
};
重连机制设计:指数退避与断点续传
连接断开后立即重连是一种直接但危险的做法。如果服务端因为压力过大而主动断开连接,或者正处于发布重启过程中,所有客户端同时重连会形成惊群效应,进一步加重服务端负担,甚至导致新的崩溃。合理的策略是采用指数退避算法:第一次重连等待1秒,之后每次失败将等待时间翻倍,例如2秒、4秒、8秒,并设置一个最大等待上限(如30秒)。同时加入随机抖动,避免大量客户端在同一时间窗口内同时发起重试。抖动范围通常取等待时间的20%到50%,例如等待4秒时,实际延迟在3.2秒到4.8秒之间随机取值。这样可以有效分散重连压力。
断点续传是更高级的需求。SSE规范支持客户端在重连时携带Last-Event-ID头,告知服务端最后成功接收的事件ID。服务端可以根据这个ID从对应位置继续发送后续事件,而不是从头开始。对于AI智能体而言,流式输出的是逐步生成的文本或任务状态更新,如果重连后重新发送全部内容,前端需要处理重复数据,还可能造成状态错乱。实现断点续传要求服务端为每个SSE事件分配唯一的递增ID(在事件数据中以id:字段发送),并且能够根据ID查询或重建后续事件流。一种简单的实现是服务端将已生成的事件按顺序缓存在内存或Redis中,键为任务ID,值为事件列表。重连时读取Last-Event-ID,从该ID之后开始发送。如果缓存中不存在该ID(例如过期),则回退到全量发送或返回错误码让前端重新发起任务。
以下是一个实现了指数退避重连和Last-Event-ID传递的JavaScript客户端代码。它手动管理重试计时器,避免同时存在多个EventSource实例:
class SSEWithReconnect {
constructor(url, options = {}) {
this.url = url;
this.lastEventId = null;
this.retryCount = 0;
this.maxRetry = options.maxRetry || 5;
this.baseDelay = options.baseDelay || 1000; // 初始延迟1秒
this.maxDelay = options.maxDelay || 30000; // 最大延迟30秒
this.jitter = options.jitter || 0.3; // 抖动系数
this.eventSource = null;
this.timer = null;
this.listeners = { message: [], error: [] };
}
connect() {
// 构造URL,如果存在lastEventId则通过查询参数传递(也可依赖浏览器自动携带的Last-Event-ID头)
const url = new URL(this.url, window.location.origin);
if (this.lastEventId) {
url.searchParams.set('last_event_id', this.lastEventId);
}
this.eventSource = new EventSource(url.toString());
this.eventSource.onmessage = (event) => {
// 更新lastEventId
this.lastEventId = event.lastEventId || this.lastEventId;
// 重置重试计数,表示连接成功
this.retryCount = 0;
this.emit('message', event);
};
this.eventSource.onerror = (event) => {
this.emit('error', event);
this.eventSource.close();
this.scheduleReconnect();
};
}
scheduleReconnect() {
if (this.retryCount >= this.maxRetry) {
console.error('达到最大重试次数,停止重连');
return;
}
const delay = Math.min(
this.baseDelay * Math.pow(2, this.retryCount),
this.maxDelay
);
// 加入随机抖动,范围 [delay * (1 - jitter), delay * (1 + jitter)]
const jitteredDelay = delay * (1 - this.jitter + Math.random() * 2 * this.jitter);
this.retryCount++;
console.log(`第${this.retryCount}次重连,延迟${Math.round(jitteredDelay)}ms`);
this.timer = setTimeout(() => {
this.connect();
}, jitteredDelay);
}
emit(type, event) {
this.listeners[type].forEach(fn => fn(event));
}
on(type, fn) {
this.listeners[type].push(fn);
}
close() {
clearTimeout(this.timer);
if (this.eventSource) {
this.eventSource.close();
}
}
}
// 使用示例
const sse = new SSEWithReconnect('/api/agent/stream?task_id=123');
sse.on('message', (event) => {
// 处理消息
console.log(event.data);
});
sse.connect();
上面的代码中,lastEventId通过查询参数手动传递,实际上现代浏览器在EventSource重连时也会自动在请求头中加入Last-Event-ID,不过这个自动行为只发生在浏览器自动重连时(但浏览器不会自动重连),所以我们手动控制更加可靠。服务端需要解析这个参数或请求头,实现续传逻辑。需要注意的是,如果服务端每次SSE事件都携带了id:字段,前端可以从event.lastEventId中获取到该ID,但如果服务端未提供,则lastEventId为空,此时断点续传无法启用,只能全量重发。
生产环境实践:心跳检测与优雅降级
即使实现了重连和续传,长连接仍然可能因为空闲而被中间设备关闭。许多代理和负载均衡器的空闲超时时间可能在30秒到60秒之间,如果AI智能体的流式输出间隔恰好超过这个时间(例如模型思考时间较长),连接就会在没有任何数据的情况下被静默断开。为了防止这种情况,需要引入心跳机制。心跳可以通过两种方式实现:一是服务端定期发送SSE注释行(以冒号开头),注释行不会触发客户端的onmessage事件,但能保持连接活跃;二是服务端发送自定义事件(如event: heartbeat),客户端可以监听该事件并更新最后的活跃时间戳。推荐使用注释行,因为注释行对应用逻辑透明,不会干扰正常的消息处理。
客户端也需要具备超时检测能力:如果长时间没有收到任何数据(包括心跳),即使连接在技术层面仍然存在,也可能已经“僵尸”了。可以在客户端设置一个看门狗定时器,每次收到消息或心跳时重置计时器,超过阈值(如心跳间隔的两倍)则强制关闭当前连接并触发重连。以下是一个带心跳检测的客户端增强示例:
// 在SSEWithReconnect基础上增加心跳检测
class SSEWithHeartbeat extends SSEWithReconnect {
constructor(url, options = {}) {
super(url, options);
this.heartbeatTimeout = options.heartbeatTimeout || 10000; // 10秒心跳超时
this.watchdogTimer = null;
}
connect() {
super.connect();
this.startWatchdog();
}
startWatchdog() {
clearTimeout(this.watchdogTimer);
this.watchdogTimer = setTimeout(() => {
console.warn('心跳超时,强制断开并重连');
if (this.eventSource) {
this.eventSource.close();
}
this.scheduleReconnect();
}, this.heartbeatTimeout);
}
// 重写onmessage,收到任何消息都重置看门狗
setEventHandlers() {
this.eventSource.onmessage = (event) => {
this.lastEventId = event.lastEventId || this.lastEventId;
this.retryCount = 0;
this.emit('message', event);
this.startWatchdog(); // 重置看门狗
};
// 监听心跳事件
this.eventSource.addEventListener('heartbeat', (event) => {
console.log('收到心跳');
this.startWatchdog();
});
this.eventSource.onerror = (event) => {
this.emit('error', event);
this.eventSource.close();
clearTimeout(this.watchdogTimer);
this.scheduleReconnect();
};
}
}
服务端心跳代码片段(Node.js + Express):
app.get('/api/agent/stream', (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
});
// 发送初始事件
res.write(`data: ${JSON.stringify({type: 'start'})}\n\n`);
// 每15秒发送心跳注释
const heartbeatInterval = setInterval(() => {
res.write(': heartbeat\n\n'); // 注释行
// 或者发送自定义心跳事件
// res.write(`event: heartbeat\ndata: {}\n\n`);
}, 15000);
// 清理
req.on('close', () => {
clearInterval(heartbeatInterval);
res.end();
});
});
最后要考虑优雅降级。当SSE连续多次重连失败,或者检测到网络环境不支持长连接(例如某些代理强制断开)时,应该切换到备选方案。一种简单的降级策略是转为普通的轮询请求:客户端定期(如每3秒)发起一个HTTP请求获取任务最新状态或新增内容。轮询虽然实时性略差,但兼容性最好,几乎所有网络环境都支持。另一种方案是升级到WebSocket,它支持双向通信且内置心跳帧,但需要服务端和客户端同时替换协议。实践中常采用混合策略:优先使用SSE,失败重试超过阈值后自动降级到轮询,待网络恢复后再尝试切回SSE。降级逻辑可以封装在客户端连接管理器中,通过状态机控制。同时,服务端也应提供对应的轮询接口,返回相同的数据结构,保证前端处理逻辑一致。监控方面,记录重连次数、重连延迟分布、心跳超时频率以及降级切换次数,这些指标有助于后续优化网络配置和服务器资源。