面试知识库

10 · 实时推送#

简历原话:实时推送:多实例部署时,事件经 Redis Pub/Sub 转发到持有该 WebSocket 连接的实例,同时写入 Redis Stream,断线重连后按 last_event_id 补发,前端按事件 ID 去重,保证不丢不重。

30 秒口述版#

ShoppingX 一轮任务要跑十几到几十秒,中间的思考、工具调用、商品卡都靠 WebSocket 实时推给前端。服务化以后任务在独立 worker 里跑,浏览器的 WebSocket 却连在某个 API 实例上,两边不在一个进程。我用两条通道分工:实时这条走 Redis Pub/Sub,worker 把事件广播出去,每个 API 实例查自己的连接表,谁持有这条会话的连接谁推;可靠这条走 Redis Stream,每个会话一条流,事件先写流、拿到 Redis 生成的单调 id 再推送。断线重连时前端带上收到的最大 last_event_id,服务端先登记新连接、再补发它之后的事件。补发和直播可能重叠,前端按事件 ID 去重,这样既不丢也不重。

背景与问题#

1. 拆出 worker 之后,事件全丢了。 单进程时,连接表是进程内的一张 thread_id → WebSocket 字典,Agent 在任意深处上报事件,按当前会话查表就能推出去。任务调度改成 Redis Stream 队列 + 独立 worker(见「任务调度」那条)以后,事件在 worker 里产生,而 worker 的连接表里一条连接都没有。现象:前端发起任务后只看到一个转圈,tool_start、商品卡预览、收尾文案一条都收不到,直到超时或者自己去轮询任务状态接口才拿到结果。结果本身是对的(状态表和历史记录都正确),丢的是整条实时过程。

2. 多副本 API,连接落在哪个实例不确定。 负载均衡把 WebSocket 分到任意一个 API 实例,用户刷新页面后新连接可能换到另一个实例。所以「事件从哪来」和「连接在哪」是两个独立的随机变量,不能靠粘性会话把它们绑在一起。

3. 断线期间的事件没了。 任务和 WebSocket 是解耦的:用户刷新、切到别的对话、手机切网、反向代理空闲超时断开,后台任务都照跑。但断开这几秒里产生的事件直接推空了。现象有两种:

  • 界面停在半截:思考流少几行,商品卡没出来;
  • 更糟的是恰好断在收尾那一刻:task_result 丢了,前端一直转圈,只有手动刷新才看到结果。

4. 补回来之后又重复。 补发的历史事件和重连后的新直播事件会有重叠。前端如果直接追加,就会出现两行一样的「正在调用 item_search」、商品卡渲染两次。

约束: 事件推送是「体验层」,不能拖垮主链路。Redis 抖动时 Agent 必须照常跑完、照常计费,最坏只是用户少看几条实时事件,最终结果仍然能从任务状态和历史里拿到。

做法与取舍#

1. 两条通道分工:Pub/Sub 管实时,Stream 管补发#

  • 实时扇出走 Pub/Sub:全集群一个频道(形如 globex:agui),worker 发布,所有 API 实例订阅。
  • 可补发存档走 Stream:每个会话一条流(键名形如 shoppingx:events:{thread_id}),每条事件 XADD 进去,Redis 返回的 stream id 就是这条事件的 ID。

为什么分开: 两件事的访问模式不一样。直播要的是「此刻扇出给所有在线实例,没有订阅者就丢掉」,Pub/Sub 不持久化在这里是需求;补发要的是「按会话、按位置查一个区间」,Stream 的 XRANGE 正好就是这个。

否掉的方案:

  • 只用 Pub/Sub:断线期间的事件没地方找回来。
  • 只用 Stream 当总线:每个会话一条流时,API 实例没法高效地「订阅所有会话」——要么对上千条流轮询 XREAD,要么再加一条全局流给每个实例各开一个消费位置,等于用 Stream 重新实现了一遍 Pub/Sub,还多一份全局持久化。
  • 用 Kafka:项目里已经有 Redis(队列、限流、计费都在用),为了几十条/轮的事件量再运维一套 Kafka 不划算;按会话范围补发在 Kafka 里也要自己建索引。

2. 路由:单频道广播 + 接收端查表过滤,不维护「会话→实例」路由表#

简历说「转发到持有该 WebSocket 连接的实例」,实现方式是:worker 把事件发到同一个频道,每个 API 实例收到后查自己的连接表,有这条会话的连接就推,没有就丢。「持有连接的那个实例」是被接收端过滤选出来的,不是发送端算出来的。

为什么这样选: 不需要一张全局路由表。路由表要在连接建立/断开时同步写 Redis,实例被 kill -9 时会留下脏路由,刷新换实例时还有「旧实例的断开回调晚于新实例的登记」的竞态,这些都要额外处理。广播方案下连接表只在进程内,天然跟着实例生死。

代价: 每个 API 实例都要收全部事件。量级是每轮任务几十条,过滤是一次字典查找,当前副本数下完全不是瓶颈。副本数上去以后的演进方式放在 Q&A 里。

否掉的方案: 按会话开频道。每次 WebSocket 连接/断开都要 SUBSCRIBE/UNSUBSCRIBE,多一套生命周期管理,换来的只是省掉一次字典查找。

3. 只在本地投不出去时才发布,并且跳过自己发的#

  • 上报事件时先尝试推给本进程的连接表,推不出去才发布到频道。worker 进程没有连接,所以它的事件全部走发布;API 进程自己手上的连接直接推成功,不再广播一遍。
  • 每个进程启动时生成一个随机 origin,发布时带在信封里,收到时比对,自己发的直接跳过。不这样做的话,多副本下同一条事件会被推 N 遍。
  • 发布是 fire-and-forget:不等 Redis 往返,不给 Agent 主循环加延迟。在途发布数设上限(1000 条),Redis「连得上但不响应」时丢新事件,不让任务堆积拖垮 worker。
  • 订阅循环自己带退避重连(断开后 2 秒重订阅)。Pub/Sub 长连接断了不会报错,只是从此收不到消息;不自己重连,一次网络抖动就会让这个副本永久没有实时事件,而且没人看得出来。
  • 订阅断开期间的频道消息是找不回来的,而浏览器那条 WebSocket 没断、不会自己重连。所以 API 实例给手上每条连接记着「最后推出去的 ID」,订阅重连成功后,对每条连接各做一次 XRANGE 补发,再恢复直播。

4. 先写 Stream,再推送:推出去的每条事件都带 ID#

上报一条事件的顺序固定为:XADD 进该会话的流 → 把返回的 stream id 回填到事件的 id 字段 → 推送(本地直推或发布到频道)。

为什么先写后推: 保证前端收到的每条直播事件都带着 ID,前端才能拿它当断点、当去重键。反过来先推后写,推出去的事件没有 ID,断线后前端不知道自己收到了哪儿。

为什么用 Redis 生成的 id 而不是自己生成: stream id 形如「毫秒时间戳-序号」,由 Redis 服务端分配,同一条流内严格单调递增,即使接管任务的是另一个 worker、两台 worker 机器时钟不一致也不影响(Redis 发现时钟回拨会沿用上一个 id 的毫秒值、只加序号)。自己用 UUID 做 ID 能去重但不能比大小,没法当补发起点;用 worker 本地自增序号,接管时序号要跨进程续上,又多一个要持久化的状态。

5. 同一会话的上报按顺序执行,所以「最大 ID」就是安全断点#

同一轮里工具是同轮并发执行的,多个协程可能同时上报事件。同一会话的「写流 + 推送」这两步按会话串行执行(进程内每个会话一把轻量锁),保证推出去的顺序和 stream id 的顺序一致。

为什么要这一步: 前端只记「收到的最大 ID」当断点。如果 id=2 的事件先推到、id=1 的事件还在路上时恰好断线,前端断点是 2,重连后补发从 2 之后开始,id=1 就永远丢了。串行上报把这个交叉窗口去掉。代价是同一会话的上报排队,但每轮只有几十条、每条是一次 Redis 往返,排队时间可以忽略;不同会话之间互不影响。

6. 重连时先登记连接,再补发#

WebSocket 握手流程:

  1. 校验 token 与会话归属(不通过直接关,不先 accept 再赶人);
  2. 把新连接登记进本实例的连接表;
  3. 回一条 ws_ready 控制帧;
  4. 如果带了 last_event_id,用 XRANGE 以排他下界 (last_event_id 到 + 查出缺口事件,逐条补发;
  5. 之后就是正常直播。

为什么先登记再补发: 登记之后产生的新事件一定能被直播推到(本实例直推,或经频道被本实例过滤命中);登记之前产生的事件一定已经在流里(先写后推)。两段覆盖了整条时间线,中间没有空档。代价是两段可能重叠:一条事件在「登记后、补发查询前」写入流,它既会被直播推一次,也会被 XRANGE 查到再补一次。重叠交给前端去重处理。

否掉的方案: 先补发再登记。补发查询和登记之间产生的事件,既不在补发结果里(查询已经结束),也推不到(连接还没登记),就丢了。这个窗口很短,但正好是用户刷新后最容易撞上的时刻。

首次连接同理(connect-first): 前端先建 WebSocket,收到 ws_ready 之后才 POST 发起任务,保证任务上报第一条事件时连接已经在表里。

7. 前端:按 ID 去重,记最大 ID,指数退避重连#

  • 维护一个已收事件 ID 集合,每条带 ID 的事件先查集合,见过就丢弃;
  • 断点 lastEventId 取收到过的最大 ID,比较时按「毫秒-序号」两段数值比较,不按字符串比较(字符串比较在位数不同时会出错);
  • 非正常断开(不是收尾后服务端主动关)就自动重连:0.8 秒起步,每次翻倍,封顶 8 秒,最多 5 次;连上即清零;5 次都失败才把这一轮标成出错;
  • 重连不重发任务请求,只带 last_event_id 重新订阅;
  • 补发期间断点不推进:重连后直播的新事件可能比补发的旧事件先到(见图 2)。如果收到新事件就把断点推到最大,补发还没收完又断一次,下次重连会从新事件之后补,中间的旧事件就漏了。所以服务端补发完发一条 replay_done 控制帧,前端在收到它之前只往 ID 集合里加、断点保持为这次重连带上去的值,收到后才把断点推进到集合里的最大 ID。

刷新页面的情况: 内存里的 ID 集合没了。前端先调「当前会话是否有任务在跑」接口,服务端从流里截出「最后一个 session_created 起」的事件,也就是当前这一轮,连同用户提问一并返回。前端据此重建这一轮的界面,用这批事件初始化 ID 集合和最大 ID,再带 last_event_id 连 WebSocket 续看。

8. 流的边界:长度、过期、哪些事件不进流#

  • 长度:每条流近似裁剪到约 200 条(MAXLEN ~),够覆盖一整轮任务的全部事件,不会无界增长;
  • 过期:每次写入顺带刷新 6 小时 TTL(滑动过期),任务活跃期间不过期,结束几小时后自动回收。长度只限单条流,过期才限制流的数量——每个会话一条流,不设过期会随会话数无界累积;
  • 不进流的事件:收尾文案的逐字增量(每轮几十条,每条都是到目前为止的全文)只直播不存档。断线补回几十条渐进全文没有意义,重连后 task_result 自带定稿;
  • 只存主会话的事件:前端只连主会话,其他内部上下文的事件不写流。
  • 断点被裁掉时不做部分补发:补发前先取流里最老的一条 ID,last_event_id 比它还老说明中间一段已被裁剪,服务端改发 replay_gap 控制帧,前端走「重建当前轮 / 拉历史结果」那条路,不拿一份有洞的补发冒充完整。
  • 语义重复单独处理:worker 被接管后重新执行的那一步会写出新 ID 的事件,ID 去重挡不住。事件里带「第几步、第几个工具调用」,前端思考流按这个键更新同一行,不追加新行。

9. 降级方向:Redis 出问题只少看事件,不影响任务#

  • 流的读写包一层断路器:连续 3 次失败熔断,30 秒后半开探测。熔断期间写流直接返回「无 ID」,事件照样直推,只是这段时间断线补不回来;
  • 发布失败只记日志;订阅断开自己重连;
  • 最终结果不依赖事件通道:task_result 丢了,前端也能从任务状态接口和历史记录拿到最终清单。

流程图 / 架构图#

图 1:多实例下一条事件怎么到浏览器#

flowchart LR
    subgraph W["Worker 进程 x N"]
        A["Agent 主循环上报事件"] --> L["同会话串行锁"]
        L --> X["XADD 写会话流<br/>拿到 stream id 回填"]
        X --> T{"本进程连接表<br/>有这条会话吗"}
        T -- 有 --> D1["直接推送"]
        T -- 没有 --> P["发布到频道<br/>信封带 origin"]
    end

    subgraph R["Redis"]
        S[("会话流<br/>shoppingx:events:thread_id<br/>MAXLEN约200 TTL 6h")]
        C(("Pub/Sub 频道<br/>globex:agui"))
    end

    X -.写入.-> S
    P --> C

    subgraph API["API 实例 x M"]
        C --> F1["实例 A 订阅<br/>origin 是自己则跳过"]
        C --> F2["实例 B 订阅<br/>origin 是自己则跳过"]
        F1 --> Q1{"A 的连接表<br/>有这条会话吗"}
        F2 --> Q2{"B 的连接表<br/>有这条会话吗"}
        Q1 -- 没有 --> Drop1["丢弃"]
        Q2 -- 有 --> WS["WebSocket 推送"]
    end

    WS --> B["浏览器<br/>按 id 去重 记最大 id"]
    S -.断线重连 XRANGE 补发.-> Q2

图 2:断线重连的时序#

sequenceDiagram
    participant FE as 浏览器
    participant A as API 实例 A
    participant B as API 实例 B
    participant R as Redis 流与频道
    participant W as Worker

    W->>R: XADD e1 然后发布 e1
    R->>A: 频道消息 e1
    A->>FE: 推送 e1 前端记 last=e1
    Note over FE,A: 网络断开 A 注销连接
    W->>R: XADD e2 然后发布 e2
    R->>A: 频道消息 e2
    Note over A: 连接表里已没有这条会话 丢弃
    FE->>B: 重连 携带 last_event_id=e1
    B->>B: 校验归属 登记连接 回 ws_ready
    W->>R: XADD e3 然后发布 e3
    R->>B: 频道消息 e3
    B->>FE: 直播推送 e3
    B->>R: XRANGE 从 e1 之后到最新
    R-->>B: 返回 e2 e3
    B->>FE: 补发 e2 e3
    B->>FE: replay_done 控制帧
    Note over FE: e3 已在集合中 丢弃 断点此刻才推进到 e3

图 2 里 e3 就是「先登记、再补发」带来的重叠:它在登记之后写入,所以既被直播推到,又被补发查到,前端靠 ID 集合丢掉第二份。直播的 e3 比补发的 e2 先到,所以补发结束前前端断点一直停在 e1,收到 replay_done 才推进到 e3;这时如果又断线,重连带的仍是 e1,e2 不会被跳过。

数字怎么来的#

这条 bullet 没有量化指标,「不丢不重」是用故障注入测出来的,口径如下。

测试台: 本机 docker compose 起 2 个 API 实例 + 2 个 worker + 1 个 Redis,前面一个 nginx 做轮询负载均衡(不开粘性会话),保证重连大概率换到另一个实例。模型调用走真实接口,15 条常用 query 循环跑。

客户端: 一个脚本化的 WebSocket 客户端,行为和前端一致:带 last_event_id 重连、按 ID 去重、记最大 ID。它在每轮任务运行中随机断开 13 次,每次断开 0.55 秒后重连;其中一部分断点故意卡在收尾前后(收到 items_preview 之后立刻断),专门打「task_result 在断线窗口里」这个最坏情况。

判据(每轮都比对):

  1. 不丢:任务结束后对该会话的流做一次全量 XRANGE,截出本轮事件的 ID 集合,客户端最终收到的 ID 集合必须与之完全相等;
  2. 不重:客户端去重后交给渲染层的事件里,同一个 ID 只出现一次;去重前的原始到达数另外记下来,用来确认「重叠确实发生过、去重确实生效了」,而不是测试没打到重叠窗口;
  3. 收尾到达:每轮都必须收到 task_result。

结果: 100 轮任务、共约 200 次断线,三条判据全部通过;约三成的重连出现了补发与直播重叠,全部被去重掉。

两条反证(证明测试本身有效):

  • 关掉补发(重连不带 last_event_id):判据 1 立刻出现缺口,断在收尾附近的那些轮收不到 task_result;
  • 关掉前端去重:判据 2 出现重复 ID,页面上能看到重复的思考行。

降级验证: 任务跑到一半 docker stop 掉 Redis,任务照常跑完并写入历史,前端能从任务状态接口拿到最终结果;Redis 恢复后 30 秒内断路器半开探测成功,新事件重新带 ID。

参数来源:

参数值为什么是这个值
单流长度约 200 条统计常用 query 一整轮的存档事件数,最长的一轮也在这个数以内
流过期6 小时滑动补发只服务「短时间内断线重连」,过期后看历史记录就够了
重连退避0.8s 起翻倍,封顶 8s,最多 5 次5 次累计约 23 秒,覆盖刷新和短暂断网;再长就当后端不可用,给用户明确报错
断路器连续 3 次失败熔断,30 秒半开Redis 没起时每条事件都等连接超时会拖慢主循环,熔断后改成立即失败
在途发布上限1000 条Redis 卡住时防止发布任务无限堆积

追问 Q&A#

Q1:为什么不开负载均衡的粘性会话,让 WebSocket 和任务落在同一台机器上? 粘性只能把「同一个浏览器」粘到同一个 API 实例,粘不住「任务在哪个 worker 上跑」——任务是从 Redis Stream 队列里被任意一个空闲 worker 领走的,worker 被 kill 以后还会换一个 worker 接管。连接和任务本来就是两个独立分配的东西,要把事件从 worker 送到 API 实例这一跳是省不掉的。另外粘性会话在实例下线时照样要重连换实例,补发机制还是得有。

Q2:所有事件广播给所有 API 实例,实例多了不就浪费了? 是的,代价是 M 个实例每个都收全量。当前每轮任务几十条事件、实例个位数,过滤就是一次字典查找,不是瓶颈。副本数上到几十个时我会改成定向投递:每个 API 实例订阅一个自己的频道 globex:agui:{instance_id},连接建立时写一条 ws:route:{thread_id} → instance_id,带 TTL、由连接心跳续期;worker 先查路由再发到对应频道,查不到或路由过期就回退到广播频道。之所以现在不做,是它引入了脏路由和刷新换实例的竞态,要额外处理,而当前规模换不来收益。

Q3:Pub/Sub 本身是会丢消息的,订阅连接断一下中间的消息就没了,你怎么说「不丢」? 「不丢」靠的是 Stream,Pub/Sub 只负责「快」。两种丢法分开看:

  • 浏览器和 API 之间断了:前端会重连,带 last_event_id 从流里补;
  • 浏览器连接没断,但 API 实例和 Redis 之间的订阅断了(Redis 重启、订阅端输出缓冲超限被踢):浏览器感知不到,不会主动重连。所以 API 实例对自己持有的每条连接记着「最后推出去的 ID」,订阅循环重连成功后,对手上所有连接各做一次 XRANGE 补发,再恢复直播。前端照样按 ID 去重。

Q4:先登记再补发会重复,为什么不想办法让它不重复,而要前端去重? 要让服务端做到「恰好一次」,得在登记和补发之间加锁,挡住这段时间的直播推送,等补发完再放行,还要处理锁期间到达的事件排队。这把锁跨了「本实例连接表」和「Redis 流」两个地方,做不成原子的。反过来「服务端允许重复 + 接收端按 ID 幂等」是分布式里的常规做法:至少一次投递 + 幂等去重,效果等价于恰好一次。ID 是现成的 stream id,前端去重是一次集合查找,成本最低。

Q5:前端的 ID 集合会不会越来越大? 集合只管当前这一轮任务,一轮存档事件在 200 条以内。新一轮开始时清空重建;刷新页面时用「当前轮」接口返回的事件重新初始化。不会跨轮累积。

Q6:断线太久,断点那条事件已经被 MAXLEN 裁掉了,或者整条流过期了,怎么办? 补发前先取流里最老的一条 ID。如果 last_event_id 比它还老,说明断点之后的一段已经被裁掉,这时补发不完整,服务端不做部分补发,而是发一条 replay_gap 控制帧,前端改走刷新页面那条路:调「当前轮」接口重建本轮;如果任务已经结束,就直接从历史记录拉最终结果。流整条过期(6 小时)时同理,任务早就结束了,结果在历史里。单流 200 条是按一整轮的事件量定的,正常断线几秒到几十秒不会碰到这个分支。

Q7:worker 被 kill 后另一个 worker 接管,接管后重新执行的那一步会再发一遍事件,这算不算「重」? 这是两种不同的重复。传输层重复是同一条事件被送了两次,ID 相同,前端 ID 去重处理。接管导致的是语义重复:同一个工具调用被执行了两次,每次都真实发生,写进流的是两个不同的 ID。这个不能靠 ID 去重。处理方式是事件里带「第几步、第几个工具调用」,前端的思考流按这个键更新同一行而不是追加,所以用户看到的是「这一步的状态从运行中变成完成」,不会多出一行。副作用层面的不重复(不重复扣费、不重复加购)是任务调度和计费那两条 bullet 的事。

Q8:为什么用 WebSocket 不用 SSE?SSE 自带 Last-Event-ID 重连。 这套补发协议其实就是照着 SSE 的 Last-Event-ID 设计的:服务端给事件编号,客户端重连带上最后一个编号。选 WebSocket 是因为链路是双向的:Agent 调 ask_user 暂停等用户回答时,回答从同一条 WebSocket 发回来;心跳也走它。换 SSE 要另开 POST 接口传回复,而且 HTTP/1.1 下浏览器对同域 SSE 连接数有 6 条上限,多开几个标签页就会卡住。

Q9:每条事件多一次 XADD,会不会拖慢推送? 一次同机房 Redis 往返在亚毫秒到 1 毫秒级,每轮几十条事件,加起来几十毫秒,分摊在十几秒的任务里。模型调用一次就是秒级,不在一个数量级。Redis 慢或挂的时候断路器会熔断,写流立即失败,推送不等它。首事件延迟的压测在「限流熔断」那条 bullet 里。

Q10:同一个会话在两个标签页打开会怎样? 每个实例的连接表一个会话只登记一条连接,同一实例上新连接会覆盖旧的。注销时按连接对象身份比对,旧连接的断开回调晚到也不会误删新连接。如果两个标签页落在不同实例上,两个实例各自持有一条连接,广播时两边都会推,两个页面都能看到。设计上假定一个会话同时只有一个活跃视图,多开是能用、但不做保证的场景。

Q11:Redis 整个挂了,这条链路会怎样? 写流熔断,事件不带 ID 照样直推;发布失败只记日志,订阅端每 2 秒重试。影响是:worker 产生的事件到不了浏览器(跨进程那一跳断了),断线也补不回来。但任务本身照常跑完、照常结算,最终清单进历史记录,前端从任务状态接口拿得到。任务队列本身也在 Redis 上,Redis 挂了新任务本来就进不来,事件通道没必要比它更高可用。

Q12:断线时 Agent 正好在问用户问题(ask_user),会卡住吗? 不会。提问本身是一条存档事件,重连补发后前端照样弹出可点选的选项。回答从 WebSocket 发回到当前连着的 API 实例,那个实例手上没有等回复的任务,就经 Redis 控制通道转发给正在跑这条任务的 worker(跨进程澄清的细节不展开)。提问有 120 秒超时,断线重连通常几秒内完成。

相关八股#

1. Redis 里做消息传递的三种结构:Pub/Sub、List、Stream 有什么区别?

  • Pub/Sub:发布即推给当前在线的订阅者,不存储;订阅者不在线、断开期间的消息直接丢;一条消息扇出给所有订阅者;没有确认机制。
  • List:LPUSH + BRPOP 做简单队列,消息持久(随 RDB/AOF),但一条消息只能被一个消费者取走,取走后没有确认,消费者崩了消息就丢;不支持多组消费、不支持按位置回看。
  • Stream(5.0+):只追加的日志,每条消息有全局有序 ID;支持按 ID 区间查(XRANGE)、阻塞读新消息(XREAD)、消费者组(XREADGROUP + XACK + 待确认列表 PEL)实现至少一次;可 MAXLEN 裁剪。
  • 关联本项目:直播要扇出、不要存储,用 Pub/Sub;补发要按会话按位置查,用 Stream;任务队列要至少一次和接管,用 Stream 消费者组。

2. Stream ID 是怎么生成的,为什么能保证单调递增?

  • 格式「毫秒时间戳-序号」,由 Redis 服务端在 XADD key * 时分配;
  • 同一毫秒内序号自增;如果服务器时钟回拨,新 ID 的毫秒部分沿用流里上一条的毫秒值,只增加序号,所以单条流内严格递增;
  • 也可以显式指定 ID,但必须大于流里最后一条,否则报错;
  • 比较两个 ID 要先比毫秒再比序号,按数值比,不能按字符串比。
  • 关联本项目:用它当事件 ID,多个 worker 写同一条会话流也天然有序;前端比较 ID 时按两段数值比较。

3. XRANGE 的区间语义和 MAXLEN ~ 近似裁剪是什么?

  • XRANGE key start end,- 表示最小、+ 表示最大,默认两端闭区间;6.2 起支持在 ID 前加 ( 表示排他,比如 (1700000000000-3 + 就是「这条之后的全部」;
  • XADD key MAXLEN ~ 200 * ...:~ 表示近似裁剪。Stream 底层是基数树,每个宏节点存一批消息,近似裁剪只删整个宏节点,不会为了精确到 200 条去拆节点,所以实际长度会略超过 200,但性能好得多;
  • 精确裁剪用 =,每次都可能拆节点,开销大。
  • 关联本项目:补发用排他下界,不重发断点那一条;每条会话流用近似裁剪保留约 200 条。

4. Pub/Sub 的慢订阅者会怎样?在集群里怎么工作?

  • Redis 给每个客户端维护输出缓冲区,Pub/Sub 客户端的默认上限是 client-output-buffer-limit pubsub 32mb 8mb 60:缓冲区超过 32MB,或持续 60 秒超过 8MB,Redis 直接断开这个订阅者,缓冲区里的消息全丢;
  • 进入订阅模式的连接只能执行订阅相关命令,所以要为订阅单独开一条连接;
  • Redis Cluster 下普通 PUBLISH 会在集群所有节点间广播,节点多了带宽放大;7.0 引入分片 Pub/Sub(SPUBLISH/SSUBSCRIBE),按频道名哈希到槽,只在一个分片内传播。
  • 关联本项目:订阅端被踢或断开时浏览器不感知,所以 API 实例订阅重连后要对手上的连接主动补发一次。

5. 消息投递的三种语义:至多一次、至少一次、恰好一次,怎么实现「恰好一次」的效果?

  • 至多一次:发了不管,可能丢,不会重(Pub/Sub);
  • 至少一次:发送方收不到确认就重发,不会丢,可能重(Stream 消费者组、Kafka 默认);
  • 恰好一次:在分布式网络里端到端做不到,工程上用「至少一次 + 接收端幂等」达到同样效果,常见做法是消息带唯一 ID,接收端记录已处理的 ID 或用业务唯一键去重;Kafka 的事务/幂等生产者也只保证在 Kafka 内部。
  • 关联本项目:重连时「先登记后补发」必然产生重叠(至少一次),前端用 stream id 去重(幂等),合起来就是「不丢不重」。

6. 多实例部署时 WebSocket 连接的常见方案有哪些?

  • 粘性会话:负载均衡按 IP 或 Cookie 把同一个客户端固定到同一实例,只解决「同一客户端的多次请求落同一台」,解决不了「消息生产者在别的机器」;
  • 消息总线广播:生产者发到 Redis Pub/Sub / Kafka / MQ,所有网关实例订阅,各自查本地连接表过滤;实现简单,代价是每个实例收全量;
  • 路由表定向投递:连接建立时把「用户→网关实例」登记到 Redis 或注册中心,生产者按路由把消息发到对应实例的专属队列;省带宽,但要处理路由过期、实例宕机后的脏路由;
  • 大规模 IM 系统一般是「长连接网关层 + 路由表 + 消息存储」的组合,客户端带序号拉取补齐。
  • 关联本项目:当前用广播 + 本地过滤,规模上去后切到路由表定向投递。

7. WebSocket 怎么检测断线?为什么需要心跳?

  • TCP 连接在对端异常消失(断网、进程被 kill、NAT 表项过期)时不会立即通知,可能长时间处于半开状态,只有下一次写失败才发现;
  • 中间的反向代理有空闲超时,比如 nginx 的 proxy_read_timeout 默认 60 秒,没有数据流动就会主动断开;
  • 所以要应用层心跳:客户端定期发 ping,服务端回 pong(WebSocket 协议本身也有 Ping/Pong 控制帧),超时没收到就认为断线并重连;
  • 重连要用指数退避加上限,避免服务端重启时所有客户端同一时刻涌进来。
  • 关联本项目:前端定期发 ping 保活,断开后 0.8 秒起翻倍、封顶 8 秒、最多 5 次重连。

8. SSE 的 Last-Event-ID 机制是什么?和 WebSocket 怎么选?

  • SSE(Server-Sent Events)是基于 HTTP 的单向推送,服务端每条消息可以带 id: 字段;浏览器的 EventSource 断线后会自动重连,并在请求头里带上 Last-Event-ID,服务端据此从断点续发;retry: 字段可指定重连间隔;
  • 优点:纯 HTTP、自动重连、断点续传是协议内建的;缺点:只能服务端到客户端单向,HTTP/1.1 下浏览器同域连接数上限 6 条;
  • WebSocket:全双工、二进制和文本都支持,但重连和续传要自己做。
  • 关联本项目:需要双向(ask_user 的回答要回传),选了 WebSocket,续传协议照搬了 SSE 的 Last-Event-ID 思路。