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 的主要特征如下:

  1. 双向通信:服务端可以推送,客户端也可以主动发送命令。
  2. 消息边界明确:一个 message 事件对应一个 WebSocket message,不需要自己按换行拆包。
  3. 协议本身不保证业务消息持久化:连接断开后,已经发送但未被客户端正确处理的业务消息是否可恢复,取决于应用协议。
  4. 浏览器不能直接调用 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 恢复 需要应用协议实现 idLast-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);
}

形式化地说,对同一个实体,若服务端版本满足:

v1<v2<<vnv_1 < v_2 < \dots < v_n

客户端应用事件时要求:

apply(ei)    ei.version>localVersion\text{apply}(e_i) \iff e_i.version > localVersion

这样,重复事件不会改变状态,乱序旧事件也不会覆盖新状态。

但这个规则有一个重要边界:如果当前版本是 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 组件的职责不是“永久拥有一个连接”,而是:

  1. 在需要时订阅外部状态;
  2. 组件卸载时取消订阅;
  3. 连接管理器独立处理连接、重连和事件归并。

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 次握手;
  • 一万个客户端每秒产生十万次连接请求;
  • 服务端恢复时会遇到重连洪峰;
  • 代理、负载均衡和认证服务可能先于应用本身被压垮。

因此客户端通常使用指数退避:

dn=min(dmax,d0×2n)d_n = \min(d_{\max}, d_0 \times 2^n)

其中:

  • d0d_0 是初始延迟;
  • nn 是连续失败次数;
  • dmaxd_{\max} 是最大延迟;
  • dnd_n 是本次重连等待时间。

还应加入随机抖动。否则所有客户端在同一个故障时刻启动,退避时间相同,仍然会同时重连:

dn=dn+U(0,αdn)d'_n = d_n + U(0, \alpha d_n)

其中 UU 是均匀随机数,α\alpha 是抖动比例。

上面的代码采用最大 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 什么是半开连接

半开连接是指一端认为连接仍存在,另一端或中间网络设备实际上已经不可达。例如:

  1. 浏览器和服务端之间的 Wi-Fi 短暂中断;
  2. NAT 映射过期;
  3. 移动设备切换网络;
  4. 代理丢弃空闲长连接;
  5. 服务端进程失去响应但 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 的 SameSitewithCredentials 和服务端认证策略。若必须使用自定义 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 增量更新后周期性校验

大型列表通常采用混合方案:

  1. 事件到达时,增量更新受影响实体;
  2. 事件出现版本间隙时,立即拉取快照;
  3. 定期重新获取或比较服务端版本;
  4. 页面重新获得焦点时进行校验。

这种方式把低延迟和可靠恢复结合起来,但要求服务端定义清晰的版本或游标。


十一、使用 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

解决方案包括:

  1. 快照带版本,应用快照前保留并重放之后收到的事件;
  2. 服务端提供原子快照游标;
  3. 先暂停事件应用,拉取快照,再应用快照之后缓存的事件;
  4. 使用服务端支持的订阅游标协议。

十三、幂等、重复和乱序的完整算例

假设订单初始版本为 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 进程中的连接。部署多个实例后,某个实例产生的事件不会自动发送到其他实例上的客户端,必须引入共享事件总线或让所有实例订阅同一事件流。


十五、背压和消息过载

实时消息的生产速度可能高于浏览器渲染和业务处理速度:

λproducer>μconsumer\lambda_{\text{producer}} > \mu_{\text{consumer}}

其中 λ\lambda 是消息产生速率,μ\mu 是客户端处理速率。此时消息队列会持续增长,最终导致内存占用和页面卡顿。

处理方式取决于业务语义:

  • 股票报价:可以丢弃中间报价,只保留最新值;
  • 聊天消息:不能随意丢弃,应持久化并补偿;
  • 进度条:可以合并连续进度事件;
  • 审计事件:应保留完整序列,客户端只展示聚合结果。

例如对于高频位置事件:

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 的 SameSiteSecure 属性必须匹配部署场景;
  • 服务端仍然要检查用户是否有权订阅指定主题。

不要把长期有效的高权限令牌直接放入 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

“消息收到但界面没更新”至少有四种可能:

  1. JSON 解析失败;
  2. 版本判断认为它是旧事件;
  3. reducer/store 没有产生新引用;
  4. 组件订阅的查询键或 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 与生态项目官方文档重新梳理;正文与示例由 WR BLOG 编写。