上一篇《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)。二是重连请求经过负载均衡,也可能落到另一台机器上,这台机器手里没有这条连接之前收到的事件。

推送怎么到达持有连接的节点#

flowchart LR P["发布请求
目标:用户 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 规范)。

flowchart TB C["浏览器 / 客户端"] -->|"GET /stream(HTTP/2 上的 SSE)"| G["连接层 Stream Gateway
持有连接、本地扇出、每连接有界队列
不读写业务库"] 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),把同一个分片的用户固定到少数几个节点上,每个节点只订阅自己负责的分片。

节点只订阅本地有连接的通道#

每个节点只订阅本地有连接的房间和用户,本地连接注册表和总线订阅联动:

flowchart TB A["客户端连接
房间 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 断开连接。

sequenceDiagram participant B as 浏览器 participant N as 新节点 participant R as Relay 总线 participant H as 共享历史 B->>N: GET /stream,Last-Event-ID 为 X N->>N: 开始缓存这个通道的实时消息 N->>R: 订阅通道 N->>H: 读取 X 之后的事件 H-->>N: 历史事件 N->>N: 历史和缓存按位置合并,去掉重复 N-->>B: 先发补上的事件,再发实时事件

快照类通道和事件类通道#

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 秒左右发一行注释。以上来自公开资料和设计推演,多实例下的吞吐和延迟还要实测。

参考资料#