本节目标:把一条「能连上」的长连接做成「能长期可靠运行」的长连接。你会知道为什么 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_timeout | 60s | < 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.active | Gauge | 在线连接数,突降即故障 |
ws.reconnects.total | Counter | 按原因标签区分,识别风暴 |
ws.heartbeat.timeout | Counter | 假死连接数,反映网络质量 |
ws.broadcast.fanout_ms | Histogram | 单次广播耗时,识别慢消费者 |
指标命名与告警规则见 《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 与泛型组件 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。