SSE 多实例部署的跨节点路由与断线补发
目录
上一篇《SSE 还是 WebSocket:每连接内存、单机容量上限与多实例路由》在 4 核 8 GB 的压测环境里测了单机容量。2 万条 SSE 连接时,每条连接的 RSS 在 Netty、Go、Rust 上都是 35 KB 左右,Tomcat 全调优后是 148.1 KB;单机先后受限于 maxConnections 计数器、堆和扇出耗时。那篇最后一章写到「加第二台机器后,最先出问题的是路由」。这一篇接着讲多实例怎么设计,内容来自公开资料和设计推演,没有压测。
SSE 是一条一直不结束的 HTTP 响应,连接对象只存在于接受这条连接的那台机器的内存里。浏览器端的 EventSource 断线后会自己重连:服务端结束响应或者网络出错,它都会等一个重连间隔再发起请求;之前收到过带 id 字段的事件,重连请求就在 Last-Event-ID 头里带上最后一个 ID(见 WHATWG 的 Server-sent events)。部署成多实例后,这两个特性各带出一个问题。一是触发推送的请求常常落在另一台机器上:用户连在 Node A,更新却由 Node B 触发,Node B 在本地的连接列表里找不到这个用户,消息就发不出去(见 AYA 的 Scaling Pains: Why Your SSE App Fails Under Load Balancers)。二是重连请求经过负载均衡,也可能落到另一台机器上,这台机器手里没有这条连接之前收到的事件。
推送怎么到达持有连接的节点#
目标:用户 A"] --> N1["Node 1
没有用户 A 的连接"] N1 -.->|"消息怎么到达?"| N2["Node 2
没有用户 A 的连接"] N1 -.->|"消息怎么到达?"| NN["Node N
持有用户 A 的连接"]
跨节点路由分三层来做,后一层建立在前一层之上:连接层单独部署,总线按房间和用户路由,每个节点只订阅本地有连接的通道。
连接层单独部署#
连接层单独部署成有状态的推送服务,业务服务只负责产出事件、发到中继总线。Centrifugo 就是这样的独立服务:业务后端通过它的 server API 发布,Centrifugo 持有客户端连接并推给订阅者(见 Centrifugo design overview)。Pushpin 的 GRIP 协议让后端把长连接交给代理,后端和代理之间只走短请求(见 GRIP)。Mercure 是建在 SSE 上的协议,发布方用 POST 把更新交给 hub,订阅方按 SSE 规范连到 hub(见 Mercure 规范)。
持有连接、本地扇出、每连接有界队列
不读写业务库"] B["业务层 Producer
处理请求、提交数据库事务"] -->|"追加事件"| H["共享历史
Redis Stream 或 JetStream"] B -->|"发布带位置的帧"| R["Relay 总线
Redis Sharded Pub/Sub 或 NATS"] R -->|"按房间、按用户订阅"| G G -.->|"重连时按 Last-Event-ID 读取"| H
分成两层以后,连接层可以不带业务依赖、单独扩缩,上一篇用每连接内存的数据说明过这一点。另一个好处是部署分开:连接层每重启一次,上面的连接都要断开重连;单独部署之后,业务服务上线新版本不会断开这些长连接。图里的共享历史用于断线补发,在「重连落到另一台机器时的补发」一节展开。
按房间和用户路由#
最直接的做法是所有节点订阅同一个通道,任何节点发布的事件都广播给全部节点,由各节点在本地查找收件人。总线上的投递量是发布频率乘以节点数,以 100 个节点为例:
- 房间广播:100 个活跃房间各 1 Hz,总线每秒要投递 100 × 1 × 100 = 10,000 条。
- 用户单播:发给某个用户的一条消息也要投递给全部 100 个节点,其中 99 个节点反序列化之后发现本地没有这个用户,直接丢掉。
在 Redis Cluster 上用普通的 PUBLISH 做总线,还会多一层放大:消息会传到集群的每一个节点,不管订阅者连在哪个分片上(见 Redis Pub/Sub)。
要让消息只到达持有连接的节点,总线得按房间和用户路由。Redis 7.0 起的 Sharded Pub/Sub 用 SSUBSCRIBE 和 SPUBLISH,通道名按和 key 相同的算法分到 slot 上,消息只在持有这个 slot 的分片里转发(见 Redis Pub/Sub 的 Sharded Pub/Sub 一节)。每个房间一个通道:
SSUBSCRIBE room:101
SPUBLISH room:101 "..."
NATS 按 subject 投递,订阅了这个 subject 的客户端各收到一份(见 Core NATS)。新建 subject 几乎没有开销,NATS 能处理数百万个不同的 subject(见 Subjects),所以可以给每个用户一个单播 subject,比如 stream.user.{userId}:只有持有这个用户连接的节点订阅它,其余节点收不到这条消息。
用户单播不适合按哈希分到少量共享通道。按 hash(userId) % N 把用户分到 N 个通道时,只要 N 比单个节点上的用户数小得多,几乎每个通道里都有这个节点的用户,每个节点都得订阅全部 N 个通道,结果和全局广播差不多。可行的做法有两种:用 NATS 给每个用户一个 subject;或者在负载均衡上按用户 ID 做一致性哈希(比如 NGINX 的 hash $key consistent),把同一个分片的用户固定到少数几个节点上,每个节点只订阅自己负责的分片。
节点只订阅本地有连接的通道#
每个节点只订阅本地有连接的房间和用户,本地连接注册表和总线订阅联动:
房间 101"] --> Q{"本节点第一条
房间 101 的连接?"} Q -->|"是"| S["订阅 room:101"] Q -->|"否"| K["只增加本地计数"]
某个节点上出现某个房间的第一条连接时,订阅这个房间的通道;这个房间在本地的最后一条连接断开后退订。没有这个房间连接的节点收不到它的消息。从 Centrifugo 拆出来的 Go 库 Centrifuge 就是这样做的:本节点出现某个通道的第一个订阅者时才向 broker 订阅;通道在本节点空了以后,至少等 1 秒,确认仍然没有订阅者才退订,这 1 秒里有新连接进来就不退订(见 node.go 的 addSubscription 和 removeSubscription)。
重连落到另一台机器时的补发#
EventSource 重连时带上的 Last-Event-ID,是它收到的最后一个事件的 id 字段。服务端只要给每个事件写上 id,重连请求就能告诉服务端从哪里接着发。单机时,节点在内存里留一段最近的事件就够了;多实例下,重连请求经过负载均衡可能落到任何一个节点上。负载均衡按用户做一致性哈希也避免不了:NGINX 的 hash ... consistent 在增删服务器时只重新映射一小部分 key(见 ngx_http_upstream_module),但下线节点上的用户总要换到别的节点去。
补发要回答四个问题:历史放在哪里,新节点怎样把历史和实时消息接起来,哪些通道要逐条补发,一条连接订阅了多个通道时 Last-Event-ID 怎么表示。
历史放在共享存储里#
历史只留在原节点的内存里,重连落到别的节点就查不到。Centrifugo 默认的 Memory engine 只支持单节点,历史放在进程内存里,重启就没了;换成 Redis 后,历史存进 Redis Stream,多个节点通过 Redis 的 PUB/SUB 通信(见 History and recovery 和 Engines)。Mercure 的开源 hub 只有一个节点,多节点要用授权版,节点之间靠 Redis、PostgreSQL、Kafka 或 Pulsar 这类 transport 同步(见 High availability)。
事件的 id 最好直接用它在共享历史里的位置,哪个节点收到重连请求,都能按这个位置往后读。推给在线连接的帧也要带同一个位置,所以发布时先写历史、拿到位置,再发到总线。Centrifuge 的 Redis broker 把这几步放在同一段 Lua 脚本里:先给消息分配 offset,XADD 进这个通道的 Stream,再把带着 offset 和 epoch 的消息发布出去(见 broker_history_add_stream_fn.lua)。
实时总线本身不存消息,历史要放在能按位置读取的存储里:
| 存储 | 按位置读取 | 出处 | |
|---|---|---|---|
| Redis Pub/Sub、Sharded Pub/Sub | 不存,at-most-once | 不支持 | Redis Pub/Sub |
| Redis Stream | 存在 Redis 里,可以用 MAXLEN 限制长度 |
条目 ID 单调递增,XRANGE 的起点加 ( 前缀可以取某个 ID 之后的条目 |
Redis Streams |
| NATS Core | 不存,at-most-once | 不支持 | Core NATS |
| NATS JetStream | 追加到 stream,重启后还在 | 每条消息有序号,consumer 用 by_start_sequence 从指定序号开始读 |
JetStream、Stream and consumer policies |
| Kafka | 持久化在 topic 里,按配置保留 | 同一个分区按写入顺序读取,事件读过以后不删,可以反复读 | Kafka 文档 |
Centrifugo 的 NATS broker 就是因为 NATS Core 只有 at-most-once,历史和断线恢复都不可用(见 Engines)。
新节点接上实时消息的顺序#
新节点收到带 Last-Event-ID 的重连请求后,既要补历史,又要接上实时消息,两步的先后决定会不会漏。先读历史再订阅,两步之间发布的事件会漏掉;先订阅再读历史,同一个事件可能从两边各来一次。Centrifuge 的做法是先开始缓存这个通道的实时消息,再订阅,再读历史,最后把历史和缓存的消息按 offset 排序、去掉重复的,并检查中间有没有断档(见 client.go 和 helpers.go 的 MergePublications)。缓存期间进来的消息太多,或者历史的 epoch 变了,合并不上,就以 DisconnectInsufficientState 断开连接。
快照类通道和事件类通道#
Centrifugo 按通道选恢复模式(见 History and recovery)。默认的 stream 模式按顺序补发错过的全部消息,适合聊天、动态流、活动日志这类每条都有意义的事件;cache 模式只补发最新的一条,适合价格、仪表盘、在线状态这类每条都是完整状态快照的通道,通常配 history_size: 1。如果房间推送的是周期性的状态帧,就属于后一类,历史只留一条;发给某个用户的通知属于前一类。
历史只保留有限的条数和时间,客户端离线太久,要补的位置可能已经被截掉了。Centrifugo 在 offset 之外还记一个 epoch 标识这份历史,错过的消息已经过期,或者历史丢了(epoch 变了),都返回 recovered: false;文档要求应用这时从自己的数据库加载完整状态,恢复只优化常见的重连路径,替代不了业务库这个数据源。Mercure 规范给了 SSE 层面的判断办法:hub 可以因为运维原因丢弃事件,重连请求带了 Last-Event-ID 时,hub 必须在响应头里回一个 Last-Event-ID;订阅方拿它和自己请求的 ID 比较,对不上就说明丢了数据,应该重新获取这个 topic(见 Mercure 规范 的 Reconciliation 一节)。EventSource 接口上没有读取响应头的属性,要用这个办法,得换成基于 fetch 的客户端,比如 fetch-event-source 的 onopen 回调能拿到完整的 Response。
一条连接只有一个 Last-Event-ID#
按规范,每个 EventSource 只记一个 last event ID,每收到一个带 id 的事件就覆盖一次。一条连接同时收房间和用户两个通道的事件时,重连请求里只剩最后到达的那个事件的 id,另一个通道补到了哪里就不知道了。Mercure 允许一个订阅方带多个 topic 匹配条件,重连时仍然只发一个 Last-Event-ID,规范要求 hub 按这个 ID 补发之后的所有事件(SHOULD,见 Mercure 规范)。Centrifugo 的单向 SSE 传输干脆不支持 Last-Event-ID,文档给的理由就是一条连接可以订阅多个通道(见 Unidirectional SSE);它自己的 SDK 按通道分别记 epoch 和 offset,重连时带回去(见 History and recovery)。
设计通道时,要么让一条 SSE 连接只对应一份历史,要么让同一个 ID 能在所有订阅的通道里定位,两样都做不到就只能像 Centrifugo 的单向 SSE 那样放弃 Last-Event-ID。拆成每个通道一条连接也有代价,不走 HTTP/2 时很快会碰到浏览器每个域名 6 条连接的上限(见「浏览器连接数和中间代理」一节)。
队列满了就断开慢连接#
客户端读得慢(弱网、丢包)时,服务端写这条连接会卡住。如果扇出直接在总线的消费回调里同步写 socket,一条慢连接就会拖住消费线程,同一个节点上其他连接的消息也跟着延迟。NATS 服务端的处理办法是:客户端读得太慢,服务端在每个客户端的写超时内写不完,就放弃这个客户端,关闭整条连接(见 Slow Consumers)。
连接层可以用同样的办法:每条连接一个固定容量的队列,总线消费线程只用非阻塞的 offer() 往里放,放不进去就关闭这条连接,消费线程不会被一条慢连接挂住。示意代码:
public boolean enqueue(Frame frame) {
boolean queued = this.queue.offer(frame);
if (!queued) {
// 队列已满:判定为慢连接,关闭它,让客户端重连
this.close();
return false;
}
return true;
}
上一篇的压测 harness 就是这样做的:每条连接一个容量 256 的 ArrayBlockingQueue,队列满就关闭连接。Centrifugo 也给每条连接一个独立的消息队列,默认上限 1 MB(见 Configuration 的 client.queue_max_size)。容量越小,慢连接越早被断开;容量越大,每条连接占的内存越多,要按帧大小乘以队列长度乘以连接数来估算。
断开慢连接的代价不大,前提是补发已经做好:浏览器过一个重连间隔就会带着 Last-Event-ID 回来,从共享历史里补上断开期间的事件。
重连间隔和随机抖动#
EventSource 怎么重连,上一篇的「重连:SSE 由浏览器负责,WebSocket 要自己写」一节讲过:等待时间由 retry 字段设置,规范允许浏览器加退避但不强制,也没提抖动。放到多实例下,连接层节点滚动升级或崩溃时,上面的连接会同时断开,大量客户端按相近的间隔同时重连,新建连接会集中打到负载均衡和剩下的节点上。
上一篇也提到了规范里的 fail the connection 分支:状态码不是 200 或者 Content-Type 不对,readyState 变成 CLOSED,浏览器不再重连;规范引言里说可以用 204 No Content 让客户端停止重连,依据的就是这个分支。多实例下要多防一种情况:滚动升级期间负载均衡给重连请求回 502 或 503,原生 EventSource 同样会就此停下,页面代码要在 error 事件里检查 readyState,是 CLOSED 就重新创建连接。新建的 EventSource 的 last event ID 从空字符串开始,构造函数也只有 withCredentials 一个选项,设不了请求头,所以页面要从 message 事件的 lastEventId 属性里记下最后一个 ID,放进 URL 参数交给服务端;Mercure 的 hub 就同时接受 Last-Event-ID 头和 lastEventID 查询参数(见 Mercure 规范)。
给重试间隔加随机抖动,可以把同时发生的重试错开。AWS Architecture Blog 的 Exponential Backoff And Jitter 比较了几种抖动算法,其中 Full Jitter 在 0 和退避上限之间取随机数:sleep = random_between(0, min(cap, base * 2 ** attempt))。落到 SSE 上有两种做法。一种是服务端在每条连接的第一个事件里下发带随机量的 retry:,原生 EventSource 就按这个间隔重连。另一种是客户端改用 fetch-event-source 这类基于 fetch 的库,在 onerror 里返回自己算的重试间隔,库的 README 写明连接断开或出错时重试策略完全由调用方控制,重试时它会自动带上 last-event-id 头。
用这个库要注意两点(见 src/fetch.ts):传了自定义的 onopen,默认的 Content-Type 检查就不再执行,响应要自己校验;服务端正常结束响应时,它只调用 onclose,不会自动重连,要在 onclose 里抛出异常,交给 onerror 决定是否重试。下面的示意代码沿用 README 示例的错误分类(4xx 里除了 429 都不再重连),另外把 204 也当作停止重连的信号:
// 示意:Full Jitter,退避上限 30 秒
function getReconnectDelay(attempt) {
return Math.random() * Math.min(30000, 1000 * 2 ** attempt);
}
class FatalError extends Error {}
let attempt = 0;
fetchEventSource('/stream', {
async onopen(response) {
const type = response.headers.get('content-type') || '';
if (response.ok && type.startsWith('text/event-stream')) {
attempt = 0;
return;
}
if (response.status === 204 || (response.status >= 400 && response.status < 500 && response.status !== 429)) {
throw new FatalError(`stop reconnecting: ${response.status}`);
}
throw new Error(`retry: ${response.status}`); // 502、503 等交给 onerror 重试
},
onclose() {
throw new Error('stream closed'); // 节点下线时服务端正常结束响应,也要重连
},
onerror(err) {
if (err instanceof FatalError) {
throw err; // 重新抛出,fetchEventSource 就不再重试
}
return getReconnectDelay(attempt++); // 下一次重试前等待的毫秒数
},
});
浏览器连接数和中间代理#
连接层对外要启用 HTTP/2。不走 HTTP/2 时,浏览器对同一个域名最多只开 6 条 SSE 连接,多开几个标签页就不够用;走 HTTP/2 时,并发流数由浏览器和服务端协商,默认 100(见 MDN 的 EventSource)。WHATWG 规范还给了一个办法:同一个站点的多个页面通过 shared worker 共用一个 EventSource。
SSE 的响应要边产生边送到浏览器,中间的反向代理不能把它缓冲起来。NGINX 默认开着 proxy_buffering,SSE 的 location 要把它关掉,或者由上游在响应头里带 X-Accel-Buffering: no,只对这个响应关闭缓冲(见 NGINX 的 ngx_http_proxy_module)。代理还会断开长时间没有数据的连接:NGINX 的 proxy_read_timeout 默认 60 秒,两次读取之间超过这个时间,连接就被关闭。规范建议每 15 秒左右发一行以 : 开头的注释,防止旧式代理在短超时后断开连接,浏览器解析时会忽略这些注释行。上一篇的压测 harness 也是 15 秒没有帧就发一次心跳。
小结#
SSE 部署成多实例后,协议的两个特性各要求一套设计。连接只在一台机器上,所以连接层单独部署,业务层把事件发到 Redis Sharded Pub/Sub 或 NATS,节点只订阅本地有连接的房间和用户。浏览器会带着 Last-Event-ID 自动重连,重连又可能落到另一台机器,所以事件的 id 用共享历史里的位置(Redis Stream 的条目 ID 或 JetStream 的序号),新节点先缓存实时消息、再订阅、再读历史,按位置合并;快照类通道只留最新一条,历史不够时回业务库拉全量;一条连接订阅多个通道时,一个 Last-Event-ID 记不下多个位置。补发做好以后,慢连接可以直接断开。重连间隔要加随机抖动,非 200 的响应会让原生 EventSource 停止重连;代理要关掉缓冲,并且每 15 秒左右发一行注释。以上来自公开资料和设计推演,多实例下的吞吐和延迟还要实测。
参考资料#
- HTML Living Standard: Server-sent events(WHATWG):
id和Last-Event-ID、重连间隔和retry:、非 200 响应放弃连接、注释行和 shared worker 的建议 - EventSource(MDN):不走 HTTP/2 时每个浏览器每个域名 6 条连接的限制
- Mercure 规范、High availability:建在 SSE 上的 hub、
Last-Event-ID补发和丢数据检测、多节点同步 - Centrifugo design overview、Engines、History and recovery、Unidirectional SSE、Configuration:独立的连接服务、各 engine 的历史存储、两种恢复模式、单向 SSE 不支持
Last-Event-ID、每连接消息队列 - centrifugal/centrifuge(commit
2c3e514):node.go的按需订阅和延迟退订,Redis broker 的发布脚本,client.go和internal/recovery/helpers.go的历史与实时消息合并 - GRIP(Pushpin,Fastly 维护):把长连接交给代理的协议
- Redis Pub/Sub、Redis Streams:全局和 Sharded Pub/Sub、Stream 的条目 ID 和范围读取
- Core NATS、Subjects、JetStream、Stream and consumer policies、Slow Consumers:投递语义、subject 的开销、按序号读取、慢消费者处理
- Kafka 文档:持久化的 topic 和分区内的顺序
- ngx_http_proxy_module、ngx_http_upstream_module:
proxy_buffering、X-Accel-Buffering、proxy_read_timeout、hash ... consistent - Exponential Backoff And Jitter(Marc Brooker,AWS Architecture Blog,2015-03-04)
- fetch-event-source:基于
fetch、可自定义重试策略的 SSE 客户端 - Scaling Pains: Why Your SSE App Fails Under Load Balancers(AYA,2026-02-10):负载均衡后面的 SSE 路由问题
- SSE 还是 WebSocket:每连接内存、单机容量上限与多实例路由:单机容量的压测数据,以及
EventSource重连机制的细节