React 基础体系 · 第 48/70 篇。示例以 React 19、现代 TypeScript 和主流框架能力为基础;客户端与服务端边界会明确说明。
React 实时数据:WebSocket、SSE、重连、心跳和缓存同步
实时数据不是“收到消息后调用一次 setState”这么简单。一个完整的实时功能至少要处理以下问题:
- 服务端如何把数据推送到浏览器;
- WebSocket 与 SSE 分别提供什么通信语义;
- 连接断开后客户端如何重连,如何避免重连风暴;
- 如何区分“连接还在”与“连接实际上已经失效”;
- 消息是否可能重复、乱序或丢失;
- 组件卸载、路由切换和 React 开发模式重复挂载时,如何避免重复连接;
- 实时事件如何更新或失效已有缓存;
- 断线期间发生的数据变化如何补齐。
本文以 React 19、现代 TypeScript 和浏览器原生 API 为基础。示例中的 React 代码运行在客户端,Node.js 代码运行在服务端;两者不能混用。
一、先区分三类“实时”
“实时”在工程上通常混合了三个不同目标。
低延迟表示消息从服务端产生到客户端可见的时间较短。它不代表消息一定不丢失。
持续推送表示服务端可以主动向客户端发送数据,而不需要客户端不断轮询。WebSocket 和 SSE 都具备这一能力。
状态最终一致表示客户端最终会收敛到服务端状态。即使连接中断,只要客户端能通过重新拉取或增量补偿恢复状态,就可以实现最终一致。
因此:
低延迟 ≠ 不丢消息
持续推送 ≠ 自动重连
连接未关闭 ≠ 网络路径仍然可用
收到事件 ≠ 本地缓存已经正确同步
一个可靠设计通常把实时通道当作“状态变化通知或增量传输通道”,把 HTTP 查询或快照接口当作恢复和校验手段。
二、WebSocket 与 SSE 的通信模型
2.1 WebSocket:双向、长连接、消息边界明确
WebSocket 在 HTTP 握手成功后,把连接升级为一个双向消息通道:
客户端 ──HTTP Upgrade──> 服务端
客户端 <──101 Switching Protocols── 服务端
客户端 ──消息──────────> 服务端
客户端 <──消息────────── 服务端
浏览器 API 是:
const socket = new WebSocket("wss://example.com/realtime");
socket.addEventListener("open", () => {
socket.send(JSON.stringify({ type: "subscribe", topic: "orders" }));
});
socket.addEventListener("message", (event) => {
console.log(event.data);
});
socket.addEventListener("close", (event) => {
console.log(event.code, event.reason);
});
socket.addEventListener("error", (event) => {
console.error("WebSocket error", event);
});
WebSocket 的主要特征如下:
- 双向通信:服务端可以推送,客户端也可以主动发送命令。
- 消息边界明确:一个
message事件对应一个 WebSocket message,不需要自己按换行拆包。 - 协议本身不保证业务消息持久化:连接断开后,已经发送但未被客户端正确处理的业务消息是否可恢复,取决于应用协议。
- 浏览器不能直接调用 WebSocket 控制帧 ping/pong API:服务端可以发送协议级 ping,浏览器会自动处理 pong,但 JavaScript 通常不能直接观察或主动发送控制帧。因此浏览器端经常使用应用级心跳消息。
WebSocket 适合:
- 聊天和协同编辑;
- 双向交易或控制面板;
- 需要客户端频繁发送操作的应用;
- 服务端需要根据客户端指令即时响应的协议。
WebSocket 不等于“可靠消息队列”。如果业务要求断线后不丢事件,必须在消息中加入事件 ID、序号、版本号或游标,并提供补偿机制。
2.2 SSE:单向、基于 HTTP、事件流
SSE,即 Server-Sent Events,是服务端向浏览器单向推送事件的机制。浏览器通过 EventSource 建立连接:
const source = new EventSource("/api/events");
source.onmessage = (event) => {
console.log(JSON.parse(event.data));
};
source.onerror = () => {
console.log("SSE 连接出现错误");
};
服务端响应的媒体类型必须是:
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
事件流使用文本格式,每个事件以空行结束:
id: 101
event: order.updated
data: {"id":"o-1","status":"paid"}
SSE 常用字段:
event:事件类型;没有时,浏览器触发默认的message事件。data:事件数据,可以出现多行;客户端会将多行data拼接。id:事件 ID。浏览器会记录最近一次收到的 ID,自动重连时通常通过Last-Event-ID请求头告知服务端。retry:建议的重连等待时间,单位为毫秒。
SSE 的核心限制是通信方向单一:
浏览器 ──HTTP 请求──> 服务端
浏览器 <──事件流────── 服务端
客户端不能通过同一条 SSE 连接向服务端发送业务消息。如果需要发送操作,通常另用 fetch:
await fetch("/api/orders/o-1/cancel", {
method: "POST",
});
SSE 适合:
- 通知、订单状态、任务进度;
- 服务端向多个浏览器广播事件;
- 客户端不需要通过长连接发送高频命令;
- 希望复用 HTTP 认证、代理和基础设施。
SSE 使用 UTF-8 文本事件流,二进制数据和高频双向通信不是它的优势。实际部署还要确认反向代理不会缓冲响应,否则事件可能在代理缓冲区积累后才一次性到达浏览器。
2.3 两者的差异
| 维度 | WebSocket | SSE |
|---|---|---|
| 通信方向 | 双向 | 服务端到客户端 |
| 浏览器 API | WebSocket |
EventSource |
| 数据形式 | 文本或二进制消息 | UTF-8 文本事件流 |
| 浏览器自动重连 | 没有 | 有 |
| 事件 ID 恢复 | 需要应用协议实现 | 有 id 和 Last-Event-ID 约定 |
| 应用级心跳 | 通常需要自行设计 | 服务端可发送注释或数据事件 |
| 认证方式 | Cookie、握手 URL 等 | Cookie;原生 EventSource 不支持自定义请求头 |
| 典型场景 | 双向实时交互 | 单向通知和状态推送 |
“支持自动重连”不代表 SSE 能恢复所有消息。浏览器会尝试重新建立连接,但能否补齐断线期间的事件,取决于服务端是否保存事件以及是否正确处理 Last-Event-ID。
三、实时系统必须有事件协议
直接发送任意 JSON:
{ "status": "paid" }
很快会遇到无法判断事件新旧、重复和缺失的问题。更完整的事件至少需要包含:
type EntityEvent<T> = {
id: string;
topic: string;
type: "created" | "updated" | "deleted";
entityId: string;
version: number;
occurredAt: string;
payload: T;
};
各字段的作用不同:
id:事件唯一 ID,用于去重、审计和断点续传。topic:事件所属业务流,例如orders。type:事件动作。entityId:受影响的实体。version:该实体或聚合的单调递增版本。occurredAt:服务端产生时间,不能用客户端接收时间代替。payload:事件内容。
3.1 为什么需要版本号
假设服务端依次产生两个订单事件:
version=8, status="paid"
version=9, status="shipped"
由于网络、代理或客户端调度,客户端可能先收到版本 9,再收到版本 8。如果客户端无条件覆盖:
当前状态:shipped
收到旧事件:paid
最终状态:paid // 错误
正确规则是:
if (incoming.version > current.version) {
apply(incoming);
} else {
ignoreAsOldOrDuplicate(incoming);
}
形式化地说,对同一个实体,若服务端版本满足:
客户端应用事件时要求:
这样,重复事件不会改变状态,乱序旧事件也不会覆盖新状态。
但这个规则有一个重要边界:如果当前版本是 7,收到版本 9,不能简单认为版本 8 不重要。版本 8 可能改变了另一个字段,或者事件本身是非幂等操作。因此还需要检测间隙:
if (incoming.version === current.version + 1) {
apply(incoming);
} else if (incoming.version <= current.version) {
ignore(incoming);
} else {
markGapAndResync();
}
收到 version=9 而本地为 7 时,客户端知道至少缺失了一个版本,应通过补偿接口或完整快照恢复,而不是继续盲目应用。
四、消息传输、状态同步和缓存不是同一层
实时系统常见的分层如下:
flowchart LR
A[服务端数据库] --> B[事件产生器]
B --> C[WebSocket 或 SSE]
C --> D[客户端事件归并器]
D --> E[内存状态]
E --> F[React 外部 Store]
F --> G[组件渲染]
D --> H[查询缓存失效或更新]
H --> I[HTTP 重新获取快照]
每层解决不同问题:
- 数据库保存权威状态;
- 事件产生器描述状态发生了什么变化;
- WebSocket/SSE传输事件;
- 事件归并器处理去重、版本、乱序和缺口;
- 缓存层保存查询结果;
- React只负责根据状态渲染界面。
一个常见错误是把 WebSocket 连接对象放入 React 组件状态,然后让多个组件分别创建连接:
function OrderList() {
const [socket] = useState(() => new WebSocket("/orders"));
// ...
}
function NotificationBell() {
const [socket] = useState(() => new WebSocket("/orders"));
// ...
}
这会产生多个连接、重复订阅和不一致的本地状态。连接通常应提升到应用级服务、Context、Redux store 或专门的外部 store 中,组件只订阅结果。
五、React 中管理长连接的生命周期
React 组件的职责不是“永久拥有一个连接”,而是:
- 在需要时订阅外部状态;
- 组件卸载时取消订阅;
- 连接管理器独立处理连接、重连和事件归并。
useSyncExternalStore 适合把外部实时状态安全地接入 React。它的关键约束是:getSnapshot() 返回的快照在数据不变时应保持引用稳定。
下面是一个最小的客户端实时 store。该示例使用 WebSocket,并实现:
- 单一连接;
- 指数退避和随机抖动;
- 应用级心跳;
version去重和乱序保护;- React 订阅;
- 页面隐藏时降低连接活动。
// realtime-store.ts
import { useSyncExternalStore } from "react";
type Order = {
id: string;
status: "pending" | "paid" | "shipped" | "cancelled";
version: number;
};
type ServerEvent = {
id: string;
type: "order.updated" | "order.deleted";
entityId: string;
version: number;
payload?: Order;
};
type Snapshot = {
connected: boolean;
orders: ReadonlyMap<string, Order>;
lastEventId: string | null;
error: string | null;
};
const initialSnapshot: Snapshot = {
connected: false,
orders: new Map(),
lastEventId: null,
error: null,
};
let snapshot = initialSnapshot;
const listeners = new Set<() => void>();
let socket: WebSocket | null = null;
let reconnectTimer: number | undefined;
let heartbeatTimer: number | undefined;
let reconnectAttempt = 0;
let manuallyStopped = false;
function emit(next: Snapshot) {
snapshot = next;
for (const listener of listeners) listener();
}
function updateSnapshot(patch: Partial<Snapshot>) {
emit({ ...snapshot, ...patch });
}
function clearHeartbeat() {
if (heartbeatTimer !== undefined) {
window.clearInterval(heartbeatTimer);
heartbeatTimer = undefined;
}
}
function startHeartbeat(current: WebSocket) {
clearHeartbeat();
heartbeatTimer = window.setInterval(() => {
if (current.readyState !== WebSocket.OPEN) return;
// 这是应用级消息,不是 WebSocket 控制帧 ping。
current.send(JSON.stringify({ type: "ping", sentAt: Date.now() }));
}, 20_000);
}
function applyEvent(event: ServerEvent) {
const current = snapshot.orders.get(event.entityId);
// 旧版本和重复事件不会覆盖新状态。
if (current && event.version <= current.version) {
return;
}
const nextOrders = new Map(snapshot.orders);
if (event.type === "order.deleted") {
nextOrders.delete(event.entityId);
} else if (event.payload) {
nextOrders.set(event.entityId, event.payload);
}
emit({
...snapshot,
orders: nextOrders,
lastEventId: event.id,
});
}
function scheduleReconnect() {
if (manuallyStopped || reconnectTimer !== undefined) return;
// 1s, 2s, 4s ... 最大 30s。
const base = Math.min(30_000, 1_000 * 2 ** reconnectAttempt);
const jitter = Math.random() * base * 0.2;
const delay = base + jitter;
reconnectAttempt += 1;
reconnectTimer = window.setTimeout(() => {
reconnectTimer = undefined;
connect();
}, delay);
}
function connect() {
if (manuallyStopped) return;
if (socket && (
socket.readyState === WebSocket.OPEN ||
socket.readyState === WebSocket.CONNECTING
)) {
return;
}
const url = new URL("/api/realtime", window.location.origin);
if (snapshot.lastEventId) {
url.searchParams.set("lastEventId", snapshot.lastEventId);
}
const current = new WebSocket(
url.toString().replace(/^http/, "ws"),
);
socket = current;
current.addEventListener("open", () => {
if (socket !== current) return;
reconnectAttempt = 0;
updateSnapshot({ connected: true, error: null });
startHeartbeat(current);
current.send(JSON.stringify({
type: "subscribe",
topic: "orders",
}));
});
current.addEventListener("message", (message) => {
if (socket !== current) return;
let data: unknown;
try {
data = JSON.parse(message.data);
} catch {
updateSnapshot({ error: "收到无法解析的实时消息" });
return;
}
if (
typeof data === "object" &&
data !== null &&
"type" in data &&
data.type === "pong"
) {
return;
}
applyEvent(data as ServerEvent);
});
current.addEventListener("error", () => {
// 浏览器通常还会随后触发 close;重连逻辑集中放在 close。
updateSnapshot({ error: "实时连接发生错误" });
});
current.addEventListener("close", () => {
if (socket !== current) return;
socket = null;
clearHeartbeat();
updateSnapshot({ connected: false });
scheduleReconnect();
});
}
function start() {
manuallyStopped = false;
connect();
}
function stop() {
manuallyStopped = true;
if (reconnectTimer !== undefined) {
window.clearTimeout(reconnectTimer);
reconnectTimer = undefined;
}
clearHeartbeat();
const current = socket;
socket = null;
if (current) {
current.close(1000, "client stopped");
}
updateSnapshot({ connected: false });
}
export const realtimeStore = {
subscribe(listener: () => void) {
listeners.add(listener);
if (listeners.size === 1) start();
return () => {
listeners.delete(listener);
if (listeners.size === 0) stop();
};
},
getSnapshot() {
return snapshot;
},
getServerSnapshot() {
return initialSnapshot;
},
};
export function useRealtimeSnapshot() {
return useSyncExternalStore(
realtimeStore.subscribe,
realtimeStore.getSnapshot,
realtimeStore.getServerSnapshot,
);
}
组件只读取快照:
import { useRealtimeSnapshot } from "./realtime-store";
export function OrderList() {
const { connected, orders, error } = useRealtimeSnapshot();
return (
<section>
<p>
状态:{connected ? "已连接" : "连接中或已断开"}
{error ? `;错误:${error}` : ""}
</p>
{[...orders.values()].map((order) => (
<div key={order.id}>
{order.id}: {order.status}(版本 {order.version})
</div>
))}
</section>
);
}
这里有一个重要的 React 语义:subscribe 负责订阅外部变化,getSnapshot 负责读取当前值。组件重新渲染不应该导致每次 render 都创建新的 WebSocket;连接创建放在 store 的生命周期中,因此多个组件可以共享一条连接。
getServerSnapshot 对服务端渲染尤其重要。服务端没有浏览器 WebSocket,也不能在渲染阶段访问 window。服务端快照应是稳定的初始值,客户端水合后再建立连接。
六、React Strict Mode 与重复连接
React 开发模式下,Strict Mode 可能执行一次额外的“挂载—清理—再挂载”流程,用来发现副作用清理问题。它不是生产环境中同样的网络行为,但会暴露以下错误:
useEffect(() => {
const socket = new WebSocket(url);
socket.onmessage = onMessage;
}, []);
如果没有返回清理函数,开发环境中可能留下旧连接。即使生产环境只有一条连接,这也说明生命周期不完整。
正确的组件级副作用应清理资源:
useEffect(() => {
const controller = new AbortController();
fetch("/api/orders", { signal: controller.signal });
return () => {
controller.abort();
};
}, []);
对于共享实时连接,不能简单地让每个组件各自关闭连接。应采用引用计数、单例 store 或 Context 管理“最后一个订阅者离开后再关闭”的策略。上面的 listeners.size 就是一个最小引用计数实现。
七、重连:不是固定延迟重试
7.1 为什么不能立即无限重连
假设客户端在服务端故障期间每 100 ms 重连一次:
- 单个客户端每秒发起 10 次握手;
- 一万个客户端每秒产生十万次连接请求;
- 服务端恢复时会遇到重连洪峰;
- 代理、负载均衡和认证服务可能先于应用本身被压垮。
因此客户端通常使用指数退避:
其中:
- 是初始延迟;
- 是连续失败次数;
- 是最大延迟;
- 是本次重连等待时间。
还应加入随机抖动。否则所有客户端在同一个故障时刻启动,退避时间相同,仍然会同时重连:
其中 是均匀随机数, 是抖动比例。
上面的代码采用最大 30 秒、最多额外 20% 随机延迟。具体数值不是协议规范,而是可根据业务和容量调整的经验参数。
7.2 何时重置失败次数
连接成功后通常把 reconnectAttempt 重置为 0。但“TCP 连接建立”不一定代表服务可用。如果连接建立后立即关闭,重置可能导致反复快速重连。
更稳妥的判定是:连接保持超过一个稳定窗口,或者完成订阅确认后,再重置退避计数。例如:
let stableTimer: number | undefined;
socket.addEventListener("open", () => {
stableTimer = window.setTimeout(() => {
reconnectAttempt = 0;
}, 10_000);
});
socket.addEventListener("close", () => {
if (stableTimer !== undefined) {
window.clearTimeout(stableTimer);
stableTimer = undefined;
}
});
这不是 WebSocket 规范要求,而是对“连接刚打开就崩溃”故障路径的保护。
7.3 哪些关闭不应该自动重连
以下情况通常不应无限重连:
- 用户主动退出登录;
- 页面明确停止实时功能;
- 服务端返回明确的鉴权失败;
- 订阅主题不存在;
- 客户端版本过旧,服务端要求升级;
- 浏览器处于离线状态。
浏览器提供了网络状态提示:
window.addEventListener("online", () => {
// 可以立即尝试一次重连
});
window.addEventListener("offline", () => {
// 暂停主动重连,等待 online 事件
});
但 navigator.onLine 只表示浏览器认为网络接口可用,不保证目标服务可达,不能作为服务健康检查的替代品。
八、心跳:检测半开连接和活性
8.1 什么是半开连接
半开连接是指一端认为连接仍存在,另一端或中间网络设备实际上已经不可达。例如:
- 浏览器和服务端之间的 Wi-Fi 短暂中断;
- NAT 映射过期;
- 移动设备切换网络;
- 代理丢弃空闲长连接;
- 服务端进程失去响应但 TCP 状态没有立即传播。
如果没有任何数据传输,客户端可能很长时间收不到 close 事件。心跳的目标是主动验证链路是否仍然可用。
8.2 WebSocket 心跳
WebSocket 有协议级 ping/pong 控制帧,但在浏览器 JavaScript API 中,应用通常不能直接发送控制帧。常见设计是:
客户端 ──{"type":"ping"}──> 服务端
客户端 <──{"type":"pong"}──── 服务端
服务端收到 ping 后回复 pong:
type ClientMessage =
| { type: "subscribe"; topic: string }
| { type: "ping"; sentAt: number };
function handleClientMessage(raw: string, send: (data: string) => void) {
const message = JSON.parse(raw) as ClientMessage;
if (message.type === "ping") {
send(JSON.stringify({
type: "pong",
sentAt: message.sentAt,
serverTime: Date.now(),
}));
}
}
仅仅 send() 成功不能证明服务端已经收到消息;客户端应记录最后一次 pong,超时后主动关闭连接,让重连逻辑接管:
let lastPongAt = Date.now();
function onPong() {
lastPongAt = Date.now();
}
const timer = window.setInterval(() => {
if (Date.now() - lastPongAt > 60_000) {
socket?.close(4000, "heartbeat timeout");
}
}, 5_000);
心跳间隔、超时阈值应考虑代理空闲超时和移动网络耗电。心跳过短会增加流量和耗电,过长则发现故障慢。它们是部署和业务相关的配置,不是固定标准值。
8.3 SSE 心跳
SSE 服务端可以发送注释行作为保活数据:
: keep-alive
以冒号开头的行不会触发业务事件,但可以让中间代理看到连接仍有数据流动。服务端也可以发送显式事件:
event: heartbeat
data: {"serverTime":1710000000000}
浏览器 EventSource 在连接关闭或发生错误时会自动尝试重连。服务端可以通过 retry 指定建议等待时间:
retry: 5000
但 SSE 自动重连不负责业务层的心跳超时判定。如果连接被某个网络设备静默吞掉,浏览器可能仍然认为连接存在;服务端注释保活和客户端超时监控仍然有价值。
九、SSE 的自动重连与事件恢复
一个简化的 SSE 服务端实现如下。这里使用 Node.js 的原生 HTTP API,服务端代码不能直接放进 React 浏览器构建产物。
// server.ts
import http from "node:http";
import { randomUUID } from "node:crypto";
type EventRecord = {
id: string;
type: string;
data: unknown;
};
const history: EventRecord[] = [];
const clients = new Set<http.ServerResponse>();
function appendEvent(type: string, data: unknown) {
const event: EventRecord = {
id: randomUUID(),
type,
data,
};
history.push(event);
// 示例只保留最近 1000 条;生产环境需根据恢复窗口设计。
if (history.length > 1000) history.shift();
const text = [
`id: ${event.id}`,
`event: ${event.type}`,
`data: ${JSON.stringify(event.data)}`,
"",
"",
].join("\n");
for (const response of clients) {
response.write(text);
}
}
const server = http.createServer((req, res) => {
if (req.url !== "/api/events") {
res.writeHead(404).end("Not found");
return;
}
res.writeHead(200, {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
});
clients.add(res);
const lastEventId = req.headers["last-event-id"];
if (typeof lastEventId === "string") {
const index = history.findIndex((event) => event.id === lastEventId);
if (index >= 0) {
for (const event of history.slice(index + 1)) {
res.write([
`id: ${event.id}`,
`event: ${event.type}`,
`data: ${JSON.stringify(event.data)}`,
"",
"",
].join("\n"));
}
} else {
// 游标太旧,无法仅靠历史补齐;通知客户端重新拉取完整快照。
res.write([
"event: resync-required",
`data: ${JSON.stringify({ reason: "cursor-expired" })}`,
"",
"",
].join("\n"));
}
}
const keepAlive = setInterval(() => {
res.write(": keep-alive\n\n");
}, 20_000);
req.on("close", () => {
clearInterval(keepAlive);
clients.delete(res);
});
});
server.listen(3000, () => {
console.log("SSE server listening on http://localhost:3000");
});
// 仅用于演示服务端产生事件。
setInterval(() => {
appendEvent("order.updated", {
id: "order-1",
status: "paid",
version: Date.now(),
});
}, 10_000);
客户端可以直接使用浏览器的自动重连:
const source = new EventSource("/api/events", {
withCredentials: true,
});
source.addEventListener("order.updated", (event) => {
const data = JSON.parse(event.data);
console.log("订单更新", data);
});
source.addEventListener("resync-required", async () => {
await fetch("/api/orders");
});
source.onerror = () => {
console.warn("SSE 连接断开,浏览器将按协议尝试重连");
};
这里的恢复有一个边界:服务端必须能根据 Last-Event-ID 找到该事件之后的历史。如果历史已经被清理,服务端无法知道客户端缺少哪些事件,只能要求客户端重新获取完整快照。
如果使用原生 EventSource,客户端不能像 fetch 那样设置任意自定义请求头。跨域场景还要同时满足 CORS、Cookie 的 SameSite、withCredentials 和服务端认证策略。若必须使用自定义 Authorization 头,通常需要使用 fetch 实现流式读取,或者改用 WebSocket;不能假设 EventSource 支持自定义 headers。
十、缓存同步:更新、失效和重新获取
“缓存同步”有三种不同策略。
10.1 直接更新缓存
如果事件包含足够完整的数据,可以直接更新实体缓存:
function applyOrderEvent(
orders: Map<string, Order>,
incoming: Order,
): Map<string, Order> {
const current = orders.get(incoming.id);
if (current && current.version >= incoming.version) {
return orders;
}
const next = new Map(orders);
next.set(incoming.id, incoming);
return next;
}
适合:
- 事件数据是完整实体;
- 版本明确;
- 列表排序、分页和权限逻辑简单;
- 客户端能正确维护派生字段。
风险是事件 payload 可能不是完整快照。例如订单事件只包含:
{
"id": "order-1",
"status": "paid"
}
如果客户端把它当作完整订单写入缓存,可能丢失金额、用户和商品信息。
10.2 使缓存失效,再重新获取
更安全的方式是把实时事件当作“缓存过期通知”:
source.addEventListener("order.updated", async (event) => {
const { id } = JSON.parse(event.data);
queryClient.invalidateQueries({
queryKey: ["order", id],
});
queryClient.invalidateQueries({
queryKey: ["orders"],
});
});
这里的 queryClient 代表具体查询缓存库的客户端;这段写法以常见 Query 类库为例,不是 React 或 SSE 自带 API。
失效策略的因果关系是:
收到“订单变化”事件
↓
不能确定本地列表是否完整
↓
标记相关查询过期
↓
下一次访问或后台刷新时获取权威快照
优点是业务正确性更容易保证,缺点是会产生额外 HTTP 请求,且事件到达后不一定立即看到新数据,取决于缓存库的刷新策略。
10.3 增量更新后周期性校验
大型列表通常采用混合方案:
- 事件到达时,增量更新受影响实体;
- 事件出现版本间隙时,立即拉取快照;
- 定期重新获取或比较服务端版本;
- 页面重新获得焦点时进行校验。
这种方式把低延迟和可靠恢复结合起来,但要求服务端定义清晰的版本或游标。
十一、使用 Redux Toolkit 管理实时事件
Redux Toolkit 的 reducer 是纯函数,适合把实时事件归并成可测试的状态。下面是一个简化 slice:
import { createSlice, PayloadAction } from "@reduxjs/toolkit";
type Order = {
id: string;
status: string;
version: number;
};
type OrdersState = {
entities: Record<string, Order>;
lastEventId: string | null;
needsResync: boolean;
};
const initialState: OrdersState = {
entities: {},
lastEventId: null,
needsResync: false,
};
type OrderUpdated = {
eventId: string;
entity: Order;
};
const ordersSlice = createSlice({
name: "orders",
initialState,
reducers: {
orderUpdated(state, action: PayloadAction<OrderUpdated>) {
const incoming = action.payload.entity;
const current = state.entities[incoming.id];
if (current && incoming.version <= current.version) {
return;
}
if (
current &&
incoming.version > current.version + 1
) {
state.needsResync = true;
return;
}
state.entities[incoming.id] = incoming;
state.lastEventId = action.payload.eventId;
},
orderDeleted(
state,
action: PayloadAction<{
eventId: string;
id: string;
version: number;
}>,
) {
const current = state.entities[action.payload.id];
if (current && action.payload.version <= current.version) {
return;
}
delete state.entities[action.payload.id];
state.lastEventId = action.payload.eventId;
},
resyncCompleted(state, action: PayloadAction<OrdersState>) {
state.entities = action.payload.entities;
state.lastEventId = action.payload.lastEventId;
state.needsResync = false;
},
},
});
export const {
orderUpdated,
orderDeleted,
resyncCompleted,
} = ordersSlice.actions;
export default ordersSlice.reducer;
Redux Toolkit 使用 Immer,使 reducer 中看似直接修改 state.entities 的代码仍然能生成不可变更新。版本判断仍然是业务协议的一部分,不是 Redux 自动提供的能力。
连接管理可以通过 middleware 或 listener middleware 启动。核心原则是:
WebSocket/SSE 负责接收
middleware 负责转换为 action
reducer 负责确定性归并
组件负责读取 selector
这样可以在不依赖浏览器连接的情况下测试 reducer:
const state1 = reducer(undefined, orderUpdated({
eventId: "e-1",
entity: { id: "o-1", status: "paid", version: 2 },
}));
const state2 = reducer(state1, orderUpdated({
eventId: "e-0",
entity: { id: "o-1", status: "pending", version: 1 },
}));
console.assert(state2.entities["o-1"].status === "paid");
如果项目使用 Redux Toolkit Query,实时事件通常有两种接入方式:
- 事件只说明资源发生变化:调用对应查询的失效机制;
- 事件包含完整资源:调用缓存更新机制,直接修改该查询的缓存数据。
无论选哪种,都必须考虑查询参数。例如:
["orders", { status: "paid", page: 1 }]
["orders", { status: "paid", page: 2 }]
一个订单状态变化可能同时影响多个分页结果、计数查询和详情查询。只更新详情缓存,不代表列表缓存已经一致。
十二、断线恢复:快照加增量
最可靠的恢复流程通常是:
sequenceDiagram
participant C as 客户端
participant T as 实时通道
participant S as 服务端
participant D as 快照接口
C->>T: 建立连接并携带 lastEventId=100
T->>S: 查询事件 100 之后的记录
alt 历史完整
S-->>T: 事件 101、102、103
T-->>C: 发送补偿事件
C->>C: 按版本归并
else 游标过期或存在缺口
S-->>T: resync-required
T-->>C: 要求重新同步
C->>D: GET /api/orders/snapshot
D-->>C: 完整快照 version=200
C->>C: 替换本地状态
C->>T: 使用新游标重新订阅
end
为什么需要“完整快照”?
假设客户端最后处理到事件 100,服务端只保留最近 50 条事件,而事件 101 到 150 已经被清理。客户端即使重新连接并携带 lastEventId=100,服务端也无法恢复缺失区间。此时继续发送最新事件会让客户端永远不知道中间发生过什么。
完整快照恢复的关键是确定快照游标。例如接口返回:
{
"orders": [
{
"id": "order-1",
"status": "shipped",
"version": 200
}
],
"cursor": "event-200"
}
客户端应将快照和游标作为一个一致性边界处理:
替换订单状态
↓
同时记录 cursor=event-200
↓
再建立实时订阅
↓
只接受 cursor 之后的事件
如果“先连接实时通道、后拉快照”却没有服务端游标约束,可能出现竞态:
客户端收到事件 201
客户端开始拉快照
快照只包含版本 200
客户端用旧快照覆盖了版本 201
解决方案包括:
- 快照带版本,应用快照前保留并重放之后收到的事件;
- 服务端提供原子快照游标;
- 先暂停事件应用,拉取快照,再应用快照之后缓存的事件;
- 使用服务端支持的订阅游标协议。
十三、幂等、重复和乱序的完整算例
假设订单初始版本为 4:
{
"id": "o-1",
"status": "pending",
"version": 4
}
客户端依次收到这些事件:
A: version=5, status=paid
B: version=7, status=shipped
C: version=5, status=paid
D: version=6, status=cancelled
逐步处理:
| 收到事件 | 本地版本 | 动作 | 结果 |
|---|---|---|---|
| A: 5 | 4 | 应用 | 版本 5,paid |
| B: 7 | 5 | 发现缺口 6 | 标记重新同步,不直接应用 |
| C: 5 | 5 | 重复旧事件 | 忽略 |
| D: 6 | 5 | 应用或等待补偿 | 取决于是否已进入 gap 状态 |
如果事件是完整实体快照,并且版本 7 包含所有字段,那么可以只使用“版本更大就覆盖”。如果事件是操作命令,例如:
{ "type": "increase-stock", "amount": 1 }
则重复执行会把库存增加两次。此时必须用事件 ID 去重:
const processedEventIds = new Set<string>();
function applyOperation(event: { id: string; amount: number }) {
if (processedEventIds.has(event.id)) return;
processedEventIds.add(event.id);
stock += event.amount;
}
内存集合只能防止当前页面生命周期内重复。若要求跨页面刷新、跨服务实例去重,去重记录应由服务端或持久化存储承担。
十四、服务端 WebSocket 示例
以下示例使用 ws 包展示服务端边界:
npm install ws
npm install -D @types/ws tsx typescript
// websocket-server.ts
import http from "node:http";
import { WebSocketServer, WebSocket } from "ws";
const server = http.createServer();
const wss = new WebSocketServer({ server, path: "/api/realtime" });
wss.on("connection", (socket, request) => {
console.log("client connected", request.url);
socket.on("message", (raw) => {
let message: unknown;
try {
message = JSON.parse(raw.toString());
} catch {
socket.close(1007, "invalid json");
return;
}
if (
typeof message === "object" &&
message !== null &&
"type" in message &&
message.type === "ping"
) {
socket.send(JSON.stringify({
type: "pong",
sentAt: Date.now(),
}));
return;
}
if (
typeof message === "object" &&
message !== null &&
"type" in message &&
message.type === "subscribe"
) {
socket.send(JSON.stringify({
type: "subscribed",
topic: "orders",
}));
}
});
socket.on("close", (code, reason) => {
console.log("client closed", code, reason.toString());
});
socket.on("error", (error) => {
console.error("socket error", error);
});
});
setInterval(() => {
const event = JSON.stringify({
id: `event-${Date.now()}`,
type: "order.updated",
entityId: "order-1",
version: Date.now(),
payload: {
id: "order-1",
status: "paid",
version: Date.now(),
},
});
for (const socket of wss.clients) {
if (socket.readyState === WebSocket.OPEN) {
socket.send(event);
}
}
}, 10_000);
server.listen(3000, () => {
console.log("WebSocket server listening on ws://localhost:3000/api/realtime");
});
这个服务器只是传输演示,不具备生产级可靠性。生产实现还需要:
- 身份认证和授权;
- 订阅主题权限检查;
- 每个连接的发送队列和背压策略;
- 事件持久化;
- 事件游标和补偿接口;
- 最大连接数和消息大小限制;
- 关闭码分类;
- 多实例之间的事件广播,例如消息代理或共享日志。
wss.clients 只包含当前 Node.js 进程中的连接。部署多个实例后,某个实例产生的事件不会自动发送到其他实例上的客户端,必须引入共享事件总线或让所有实例订阅同一事件流。
十五、背压和消息过载
实时消息的生产速度可能高于浏览器渲染和业务处理速度:
其中 是消息产生速率, 是客户端处理速率。此时消息队列会持续增长,最终导致内存占用和页面卡顿。
处理方式取决于业务语义:
- 股票报价:可以丢弃中间报价,只保留最新值;
- 聊天消息:不能随意丢弃,应持久化并补偿;
- 进度条:可以合并连续进度事件;
- 审计事件:应保留完整序列,客户端只展示聚合结果。
例如对于高频位置事件:
type LocationEvent = {
userId: string;
lat: number;
lng: number;
version: number;
};
const latestByUser = new Map<string, LocationEvent>();
function onLocation(event: LocationEvent) {
const old = latestByUser.get(event.userId);
if (!old || event.version > old.version) {
latestByUser.set(event.userId, event);
}
}
这里主动丢弃中间位置是合理的,因为最终目标是当前坐标,而不是重放每一个轨迹点。相同策略用于订单状态就可能错误,因为订单状态变化可能代表必须审计的业务事实。
十六、认证、代理和安全边界
WebSocket 使用 ws:// 或 wss://。生产 HTTPS 页面通常应使用 wss://,否则会产生混合内容或明文传输风险。
SSE 使用 https://,认证常见方式是 Cookie。跨域 Cookie 还受以下条件约束:
- 服务端必须返回正确的 CORS 允许来源;
- 使用凭据时不能使用
Access-Control-Allow-Origin: *; - Cookie 的
SameSite、Secure属性必须匹配部署场景; - 服务端仍然要检查用户是否有权订阅指定主题。
不要把长期有效的高权限令牌直接放入 WebSocket URL:
wss://example.com/realtime?token=secret
URL 可能出现在代理日志、监控日志和浏览器历史中。若架构允许,优先使用受保护的 Cookie、短期票据或在握手阶段完成受控认证。原生 EventSource 不能添加自定义 Authorization 头,这是选型时必须考虑的约束。
实时连接还要防止订阅越权。例如客户端发送:
{ "type": "subscribe", "topic": "orders:user-999" }
服务端不能仅依据字符串接受订阅,而应根据当前认证身份检查 user-999 是否属于该用户可访问范围。
十七、诊断实时功能的正确方法
17.1 先看连接层
浏览器开发者工具的 Network 面板中:
- WebSocket 请求应出现
101 Switching Protocols; - SSE 请求应保持
200且响应类型为text/event-stream; - 如果 SSE 长时间没有事件,检查代理是否缓冲;
- 如果 WebSocket 立即关闭,检查路径、升级头、TLS 和认证。
17.2 再看协议层
为每个事件打印至少这些字段:
transport=websocket
connectionId=c-123
eventId=e-456
entityId=order-1
version=17
receivedAt=...
applied=true
reason=old-version | gap | invalid | applied
“消息收到但界面没更新”至少有四种可能:
- JSON 解析失败;
- 版本判断认为它是旧事件;
- reducer/store 没有产生新引用;
- 组件订阅的查询键或 selector 不正确。
只记录“socket 收到消息”不足以定位问题,还要记录事件是否真正应用。
17.3 主动制造故障
可以在开发时进行这些验证:
- 切换浏览器离线,再恢复网络;
- 服务端重启;
- 长时间后台运行标签页后回来;
- 注入重复事件;
- 让服务端跳过一个版本;
- 让历史游标过期;
- 多次挂载和卸载实时组件;
- 同时打开多个标签页。
每种故障都应有明确预期:
断线 → 指数退避重连
重复事件 → 不改变状态
旧版本事件 → 不覆盖新状态
版本缺口 → 快照恢复
主动退出 → 不再重连
最后订阅者离开 → 关闭共享连接
十八、常见错误及其失败表现
错误一:每个组件创建一条连接
失败表现是服务端连接数快速增长、同一事件被处理多次、组件之间显示不一致。修复方式是把连接提升到应用级管理器,并让组件共享外部 store。
错误二:在 onmessage 中无条件覆盖状态
失败表现是网络乱序后,旧状态覆盖新状态。修复方式是使用服务端版本或游标,并定义旧事件、重复事件和版本缺口的处理规则。
错误三:只监听 close,不做心跳
失败表现是网络断开后页面长期显示“已连接”,直到用户刷新。修复方式是使用 WebSocket 应用级 ping/pong,或服务端定期发送 SSE 保活,并设置超时。
错误四:重连后从头加载所有数据
失败表现是断线恢复时请求量和渲染量很大,且无法判断是否漏掉中间事件。修复方式是使用事件 ID、版本或游标;游标不可恢复时再获取完整快照。
错误五:把事件流当作权威数据库
失败表现是客户端本地状态看起来正确,但由于某条事件丢失,后续状态一直错误。事件流负责传递变化,权威状态仍应由服务端快照和恢复机制保证。
错误六:SSE 服务端响应被代理缓存
失败表现是服务端日志显示不断写入,浏览器却几秒甚至几十秒才一次性收到事件。检查 Content-Type、代理缓冲设置、响应压缩和负载均衡的空闲超时。
错误七:用固定间隔无抖动重连
失败表现是服务恢复瞬间出现连接洪峰。修复方式是指数退避加随机抖动,并在服务端对连接、认证和订阅进行限流。
十九、如何选择 WebSocket 或 SSE
可以按通信需求判断,而不是按“哪个更实时”判断。
选择 SSE,当:
- 主要是服务端推送;
- 客户端操作可以通过普通 HTTP 请求发送;
- 需要文本事件、事件 ID 和浏览器自动重连;
- 希望连接模型接近 HTTP;
- 需要较简单的服务端广播。
选择 WebSocket,当:
- 客户端和服务端都需要频繁发送消息;
- 需要二进制数据;
- 需要低延迟的双向交互;
- 需要自定义订阅、确认、流控和请求响应协议;
- SSE 的自定义请求头限制不适合认证方案。
如果系统只需要每 30 秒刷新一次数据,普通轮询可能更简单、更容易观测和扩展。实时通道增加了连接生命周期、代理超时、恢复、背压和权限管理,不应为了“实时”而无条件使用长连接。
二十、最终设计模型
一个可维护的 React 实时数据系统,可以抽象为以下状态机:
IDLE
└─ subscribe → CONNECTING
CONNECTING
├─ open → CONNECTED
├─ error/close → WAITING
└─ stop → IDLE
CONNECTED
├─ event → APPLY / IGNORE / RESYNC
├─ heartbeat timeout → WAITING
├─ close → WAITING
└─ stop → IDLE
WAITING
├─ backoff elapsed → CONNECTING
├─ offline → WAITING
└─ stop → IDLE
RESYNC
├─ snapshot success → CONNECTING 或 CONNECTED
├─ snapshot failure → WAITING
└─ stop → IDLE
数据处理则遵循:
建立连接
↓
认证并订阅
↓
接收事件
↓
解析事件
↓
按 eventId 去重
↓
按 version 检查旧事件和缺口
↓
直接更新缓存,或使缓存失效
↓
发生缺口时获取快照
↓
使用游标继续接收
React 在这里不是实时协议本身,而是状态消费层。WebSocket 和 SSE 解决传输,心跳解决活性检测,重连解决连接恢复,事件 ID 和版本解决顺序与重复,快照和缓存策略解决最终一致。只有这些部分共同成立,页面上的“实时数据”才不仅是及时显示,还能在断网、重启、乱序和重复消息之后恢复到可信状态。
系列导航与关联阅读
- 系列入口:React 完整学习路线:从渲染与 Hooks 到服务端组件和生产架构
- 上一篇:Zustand 状态管理:Store、Selector、订阅、持久化和 SSR
- 下一篇:React 单元测试:纯函数、Hook、时间、网络和稳定断言
官方资料
本文依据 React 与生态项目官方文档重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论