《TypeScript编程实战》10.3 心跳、重连与广播

长连接上线前还有三块拼图没补齐:它会静默假死、断了要能自动恢复、单实例的订阅表在多副本部署下会失效。本节先讲清心跳为什么必须由应用层承担、ping/pong 与业务保活如何配合、指数退避加抖动如何避免重连风暴,再讲断线后的状态补偿与消息幂等,最后用 Redis 发布订阅把广播从单机扩展到集群,并给出背压、慢消费者与可观测性的落地方案。

本节目标:把一条「能连上」的长连接做成「能长期可靠运行」的长连接。你会知道为什么 TCP 不会告诉你对端已经死了、心跳间隔该按什么原则定、重连退避为什么必须加抖动、断线期间的消息怎么补偿,以及多实例部署时如何用 Redis 把消息广播给所有节点上的订阅者。

10.3 心跳、重连与广播

前两节解决了「消息长什么样」和「用哪种通道」。这一节处理连接的生命周期——这是长连接真正难的地方:它会在没有任何征兆的情况下静默死亡,会在网络抖动后集体重连,会在你扩容到三个副本时突然只推给三分之一的用户。

10.3.1 TCP 不会告诉你对端已经死了

最反直觉的一点:一条 TCP 连接的对端进程崩溃或网线被拔,本端的 socket 对象不会立刻报错。操作系统要等重传超时,通常是几分钟到十几分钟。这段时间里:

socket.send(payload); // 没有抛异常
console.log(socket.readyState); // 仍然是 1(OPEN)

代码认为连接还活着,用户却什么都收不到。这种「假死」状态是长连接最典型的线上故障:监控看不到错误率上升,只有用户投诉「消息不更新了」。

底层原因与排查手段可参考站内 WebSocket 实时通信架构 。解决办法只有一个:应用层自己发心跳,主动探测对端是否真的还在。

10.3.2 心跳的两层设计

心跳分两层,职责不同,都要有。

第一层是协议层 ping/pong。 WebSocket 协议内置了控制帧,浏览器会自动回 pong:

// Node 端(ws 库):服务端定时给每个连接发 ping
const HEARTBEAT_INTERVAL = 30_000;

const alive = new WeakMap<WebSocket, boolean>();

const timer = setInterval(() => {
  for (const ws of wss.clients) {
    if (alive.get(ws) === false) {
      ws.terminate(); // 上一轮没回 pong,判定为死连接
      continue;
    }
    alive.set(ws, false);
    ws.ping();
  }
}, HEARTBEAT_INTERVAL);

wss.on('connection', (ws) => {
  alive.set(ws, true);
  ws.on('pong', () => alive.set(ws, true));
  ws.on('close', () => alive.delete(ws));
});

这段代码的关键是 alive 标记的两阶段语义:发 ping 时先置 false,收到 pong 才置回 true。下一轮检查时仍是 false 的,说明整整一个周期都没回应,可以安全断开。若不这样设计,只判断「有没有 pong 事件」,就无法区分「还没回」与「已经死了」。

第二层是业务层保活。 有些中间设备(负载均衡、企业代理、云 NAT)会按「空闲时长」切断连接,而 ping/pong 属于协议控制帧,某些代理在 HTTP 升级后并不转发。所以还要有应用层消息:

type ClientMessage = { type: 'ping'; payload: { ts: number } };
type ServerMessage = { type: 'pong'; payload: { ts: number; serverTime: number } };

// 客户端:既能探活,也能顺便估算时钟偏移
setInterval(() => {
  if (socket.readyState !== WebSocket.OPEN) return;
  socket.send(JSON.stringify({ type: 'ping', payload: { ts: Date.now() } }));
}, 25_000);

时钟偏移的用途是给消息排序与延迟补偿提供依据,游戏与协同场景尤其依赖它,参见 实时系统为什么需要时钟同步 。消息结构沿用上一节的判别联合,见 《TypeScript编程实战》10.1 WebSocket 消息协议判别联合 。

间隔怎么定?取最保守的那个约束的一半:

约束来源典型值心跳上限
云负载均衡空闲超时60s< 30s
Nginx proxy_read_timeout60s< 30s
移动网络 NAT 表项30~300s< 15s
服务端连接数成本—不宜 < 10s

移动端常取 15~25 秒,桌面端 30 秒。间隔越小,探活越快,但连接数与流量成本越高——每秒一次心跳在十万连接下就是十万 QPS 的控制面流量。

10.3.3 指数退避与抖动

连接断开后立刻重连是最糟的选择。设想服务端重启,十万客户端在 1 秒内同时发起重连——服务端刚起来就被打垮,形成「重连风暴」,反复循环。

正确做法是指数退避,并叠加随机抖动:

class Reconnector {
  private attempt = 0;
  private readonly base = 500;      // 首次 500ms
  private readonly cap = 30_000;    // 上限 30s
  private readonly factor = 1.8;
  private readonly jitter = 0.3;    // ±30% 抖动

  nextDelay(): number {
    const raw = Math.min(this.cap, this.base * this.factor ** this.attempt++);
    const spread = raw * this.jitter;
    return raw - spread + Math.random() * spread * 2; // 落在 [0.7, 1.3] × raw
  }

  reset(): void {
    this.attempt = 0; // 只有「连上并稳定一段时间」后才允许归零
  }
}

三个细节决定成败:

一、抖动必须有。 没有抖动的退避只是把风暴从 1 秒推迟到 30 秒,所有客户端依然同步。加上随机抖动后,重连请求被摊平在一个区间里。

二、reset() 不能在 open 时立刻调用。 若服务端能握手但马上断开(例如鉴权通过、订阅阶段崩溃),退避永远归零,客户端会以 500ms 的间隔无限重试。正确做法是连接稳定保持 N 秒(比如 30 秒)后再归零。

三、要设上限次数并上报。 连续失败十几次后,应当停止重连、把界面切到「连接已断开,点击重试」状态。无限重连会耗干移动端电量,也让问题被掩盖。

10.3.4 断线期间的消息补偿

重连成功不等于状态一致。断线的这几分钟里,服务端可能已经推了上百条消息。三种补偿策略:

策略机制代价
全量重同步重连后重新拉一次完整状态带宽高,实现最简单
增量补偿客户端上报最后收到的序号,服务端补发需要服务端保留历史
快照 + 增量先拉快照,再补快照之后的增量最平衡,实现最复杂

增量补偿的关键是单调递增的序号,而不是时间戳——时间戳会重复、会回拨,序号不会:

type ServerMessage = { type: 'event'; payload: { seq: number; body: unknown } };

// 客户端记录收到的最大序号,重连时带上
const params = new URLSearchParams({ since: String(lastSeq) });
const socket = new WebSocket(`/ws?${params}`);

// 服务端:补发 seq > since 的消息,再切到实时推送
function onConnect(ws: WebSocket, since: number) {
  for (const msg of ringBuffer.since(since)) ws.send(encode(msg));
  liveSubscribers.add(ws);
}

ringBuffer 只需要保留最近 N 条(比如 1000 条),超过窗口的客户端就退化成全量重同步。这个「有界缓冲 + 降级」的组合是工程上最实用的方案。

补发会带来重复投递:断线前客户端可能已收到 seq=42 但还没来得及处理。因此消费端必须幂等——按 seq 去重,或让业务操作本身可重复执行。幂等设计的通用做法见 《TypeScript编程实战》9.2 重试、幂等与死信 。

10.3.5 从单机广播到集群广播

单实例时,广播就是遍历本进程的连接表:

const channels = new Map<string, Set<WebSocket>>();

function broadcast(channel: string, msg: ServerMessage): void {
  const payload = JSON.stringify(msg);
  for (const ws of channels.get(channel) ?? []) {
    if (ws.readyState === WebSocket.OPEN) ws.send(payload);
  }
}

部署到三个副本后,channels 只包含连到本进程的那部分客户端。用户 A 在副本 1、用户 B 在副本 2,两人订阅同一个房间,A 发的消息 B 永远收不到——这是长连接扩容时最经典的故障。

解法是把广播的扇出从进程内存搬到共享的中间件上。Redis 发布订阅是最轻量的选择:

import { createClient } from 'redis';

const pub = createClient({ url: env.REDIS_URL });
const sub = pub.duplicate();
await Promise.all([pub.connect(), sub.connect()]);

const NODE_ID = process.env.HOSTNAME ?? 'local';

// 发布:本进程产生的事件,投递到 Redis 频道
async function publish(channel: string, msg: ServerMessage): Promise<void> {
  await pub.publish(`rt:${channel}`, JSON.stringify({ origin: NODE_ID, msg }));
}

// 订阅:每个副本都监听所有频道,但只投递给本进程的连接
await sub.pSubscribe('rt:*', (raw, redisChannel) => {
  const { msg } = JSON.parse(raw) as { origin: string; msg: ServerMessage };
  const channel = redisChannel.slice(3); // 去掉 'rt:' 前缀
  for (const ws of channels.get(channel) ?? []) {
    if (ws.readyState === WebSocket.OPEN) ws.send(JSON.stringify(msg));
  }
});

用 pSubscribe 而不是为每个频道单独订阅,是因为频道是业务动态创建的,预先订阅无法穷举。频道的命名规范(前缀、分隔符、避免热点大频道)沿用缓存键那一套,见 《TypeScript编程实战》8.1 缓存层次与键设计 ;Redis 客户端的类型化封装见 《TypeScript编程实战》8.2 Redis 类型安全封装 。

选型上要注意 Redis pub/sub 是尽力而为:订阅者断线期间的消息直接丢失,没有持久化与 ack。需要可靠投递就得换 Streams(XADD/XREAD + 消费者组),或者直接用 Kafka 这类日志型中间件。两者的取舍见站内 Redis 发布订阅与 Streams 与 消息扇出架构设计 。

方案可靠性历史回溯适用
Redis pub/sub尽力而为无在线状态、临时通知
Redis Streams有 ack 与消费者组有界需要可靠投递的推送
Kafka强持久化、可重放完整跨系统事件流

10.3.6 背压与慢消费者

广播里最危险的是一台慢客户端:它读得慢,但服务端还在往里写,ws.bufferedAmount 一路涨到几百兆,最后把整个进程的内存拖垮。一个慢客户端能拖死整个房间,这是广播场景的头号杀手。

防护手段是给每个连接设缓冲区上限,超了就断开或降级:

const MAX_BUFFERED = 1 << 20; // 1 MiB

function safeSend(ws: WebSocket, payload: string): void {
  if (ws.bufferedAmount > MAX_BUFFERED) {
    logger.warn({ buffered: ws.bufferedAmount }, 'slow consumer, dropping');
    ws.close(1013, 'too many pending messages'); // 1013 = Try Again Later
    return;
  }
  ws.send(payload);
}

对高频场景(行情、游戏帧同步)还要进一步做合流:同一频道在 16ms 内产生多条消息时,只发最后一条状态快照,而不是逐条发送。这需要区分「状态型消息」与「事件型消息」——状态可以丢中间的、事件不能丢。相关限流与整形策略见 实时指令限流架构 。

10.3.7 消息顺序

同一连接内的消息是有序的,但经过 Redis 扇出后,来自不同副本的消息可能乱序到达。对顺序敏感的场景要做缓冲排序:给每条消息打上 seq 与 ts,接收端维护一个小窗口,窗口内的乱序消息先缓存,等齐或超时再按序交付。实现细节见 消息排序缓冲架构 。

10.3.8 可观测性与测试

长连接的问题几乎都表现为「偶发」,没有指标就只能靠猜。必须采集的四组数字:

指标类型用途
ws.connections.activeGauge在线连接数,突降即故障
ws.reconnects.totalCounter按原因标签区分,识别风暴
ws.heartbeat.timeoutCounter假死连接数,反映网络质量
ws.broadcast.fanout_msHistogram单次广播耗时,识别慢消费者

指标命名与告警规则见 《TypeScript编程实战》17.2 指标与告警 ,链路追踪(把连接 id 贯穿到每次广播)见 《TypeScript编程实战》17.1 OpenTelemetry 追踪 。

重连与心跳这类逻辑必须写测试,且要用假时钟而不是 sleep:

import { describe, expect, it, vi } from 'vitest';

describe('Reconnector', () => {
  it('退避不超过上限且带抖动', () => {
    vi.useFakeTimers();
    const r = new Reconnector();
    const delays = Array.from({ length: 12 }, () => r.nextDelay());
    expect(Math.max(...delays)).toBeLessThanOrEqual(30_000 * 1.3);
    expect(new Set(delays).size).toBeGreaterThan(1); // 抖动生效
    vi.useRealTimers();
  });
});

用 vi.useFakeTimers() 后测试从「等 30 秒」变成毫秒级完成。测试组织方式见 《TypeScript编程实战》4.1 Vitest 单元测试 。

10.3.9 与本书其它章节的衔接

心跳消息与广播消息都沿用 《TypeScript编程实战》10.1 WebSocket 消息协议判别联合 的判别联合契约;单向推送场景下 SSE 的自动重连与 Last-Event-ID 补偿见 《TypeScript编程实战》10.2 SSE 与流式响应 。进程退出时要先向所有连接发关闭帧再断开,见 《TypeScript编程实战》5.3 优雅关闭与健康检查 ;灰度发布时新旧版本连接共存的处理见 《TypeScript编程实战》18.2 数据库迁移与灰度发布 。

站内延伸阅读:可靠心跳设计 、客户端重连与会话恢复 、重连令牌与会话设计 、WebSocket 水平扩容 、WebSocket 实时测试 。

小结

长连接的可靠性由三件事共同保证。心跳解决「连接是不是还活着」——协议层 ping/pong 负责探活,业务层 ping 负责穿透中间设备,判定必须用两阶段标记而不是事件计数。重连解决「断了怎么办」——指数退避叠加随机抖动打散重连风暴,退避归零要等连接稳定之后,并且要有次数上限。补偿解决「断线期间漏了什么」——用单调序号做增量补发,客户端按序号去重保证幂等,超出缓冲窗口就降级为全量重同步。

广播则要意识到进程内存里的订阅表只覆盖本副本。把扇出交给 Redis 发布订阅即可扩展到集群,但要清楚它是尽力而为、不保留历史的;需要可靠投递就换 Streams 或 Kafka。最后别忘了给慢消费者设缓冲区上限——一个读得慢的客户端足以拖垮整个房间。

第十章到此结束。这一章从消息协议的类型建模讲到单向流式推送,再到连接的生命周期管理,覆盖了长连接从「能连上」到「能可靠运行」的完整路径。下一章转向浏览器:当这些数据要落到 React 组件里,props、泛型组件与自定义 Hook 的类型该如何设计。

阅读导航:上一节:10.2 SSE 与流式响应 · 下一节:11.1 组件 props 与泛型组件 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「typescript」更多文章

  1. 《TypeScript高级编程》11.3 类型驱动架构与团队规范
  2. 《TypeScript高级编程》11.2 渐进式迁移与严格化路径
  3. 《TypeScript高级编程》11.1 TS 版本演进与 breaking changes