← 返回文章

从 Context 到 Zustand:构建可复用的 WebSocket 实时数据层

为什么 Socket 不应该散落在组件里?

交易行情、资产变动、订单通知都要求页面在数据到达时立即更新。最直接的写法是在组件里 new WebSocket(),但应用一旦变大,很快会遇到这些问题:

  • 多个组件重复建立连接;
  • 重连、心跳、页面隐藏、网络切换逻辑散落各处;
  • 组件卸载后监听器或定时器没有清理;
  • Socket 尚未连接时发送的消息直接丢失;
  • 高频消息进入 React Context,导致整棵 Provider 子树重复渲染;
  • 行情、资产和交易页面共用一条连接,生命周期互相干扰。

我最近在一个 Next.js 13 Pages Router 交易项目中完成了 WebSocket 的 Zustand 化。项目使用 React 18、Zustand 5,最终形成了一套「连接实例在 Store 闭包中、连接状态由 Zustand 驡动、React 只负责挂载生命周期」的方案。

本文先拆解真实项目中的架构,再给出一套可以直接放进项目使用的完整代码。示例保留了现有实现的核心设计,但使用通用协议名,并修正了分析过程中发现的边界问题,不是业务源码的逐字复制。

示例只依赖 Zustand:

1
npm install zustand

当前项目的前端架构

项目仍然使用 Next.js Pages Router,核心目录职责如下:

1
2
3
4
5
6
7
8
pages/                         页面路由与页面级编排
components/ 业务组件和通用组件
components/Bootstrap/ 客户端全局能力的启动组件
store/ Zustand Store
context/ 页面会话级 Store Provider 和少量 Context
hooks/ 跨组件复用逻辑
request/ HTTP 请求封装
utils/ 消息转换、监控、跨标签页通信等工具

这里有两种 Zustand 使用方式:

  1. 用户、主题、Socket 等真正的客户端全局状态,使用 create() 创建模块级 Store;
  2. 每个交易页面独立的交易会话状态,使用 createStore() 创建 vanilla Store,再由 React Context 把 Store 实例注入当前页面。

第二种方式很重要。如果把不同交易页的价格、深度、下单状态都放进模块级单例 Store,页面切换或同时挂载多个交易区域时很容易串数据。Context 在这里不是用来承载高频业务数据,而只是负责传递一个隔离的 Store 实例。


Socket 的分层设计

当前项目没有把所有实时数据塞进一条连接,而是创建了三个相互独立的 Store:

Store 生命周期 主要职责
marketSocketStore _app.jsx 全局 市场列表、榜单、板块等公共行情
assetSocketStore 登录后全局 资产、订单通知、登录授权、资讯推送
tradeSocketStore 交易详情页 当前交易标的的价格、深度、逐笔成交

它们都由同一个 createSocketStore() 工厂创建,但可以分别配置消息过滤、消息归一化和跨标签页广播。

整体数据流如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
环境配置 / 登录状态 / 页面参数
|
v
WebSocketBootstrap
- connect / disconnect
- visibility / online
- heartbeat
|
v
createSocketStore
- WebSocket 实例
- 重连与发送队列
- 本地订阅者
- Zustand 连接状态
|
+-----+-----+
| |
v v
业务组件订阅 状态组件选择器
解析业务消息 isConnected/retryCount
|
v
React local state / 页面级 Trade Store

这套结构的关键不是“用 Zustand 保存 WebSocket 对象”,恰恰相反:

  • socketRef、订阅者数组、发送队列和定时器放在 Store 工厂的闭包中;
  • isConnectedretryCountreason 等需要驱动 UI 的值放进 Zustand;
  • 原始行情消息直接分发给订阅函数,不把每一帧消息写进全局 Store;
  • 业务组件收到消息后,只更新自己关心的 local state 或页面级 Store。

这样既能统一管理连接,又不会让每条行情触发所有 Socket 消费者重渲染。


第一步:实现跨标签页消息服务

资产消息在多个标签页之间需要同步,可以用 BroadcastChannel 做一个很薄的适配层。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
// utils/broadcastService.js
class BroadcastService {
constructor(channelName = "realtime-updates") {
this.channel = null;
this.listeners = new Set();
this.tabId = Math.random().toString(36).slice(2);

if (typeof window === "undefined" || !("BroadcastChannel" in window)) {
return;
}

this.channel = new BroadcastChannel(channelName);
this.channel.onmessage = (event) => {
const { data, originTabId } = event.data || {};
if (originTabId === this.tabId) return;
this.listeners.forEach((listener) => listener(data));
};
}

broadcast(data) {
this.channel?.postMessage({ data, originTabId: this.tabId });
}

addListener(listener) {
this.listeners.add(listener);
}

removeListener(listener) {
this.listeners.delete(listener);
}

close() {
this.channel?.close();
this.channel = null;
this.listeners.clear();
}
}

let instance;

export function getBroadcastService() {
if (!instance) instance = new BroadcastService();
return instance;
}

需要注意:这段代码只做“消息广播”,并没有选举主标签页。如果每个标签页都建立了资产 Socket,那么同一条服务端事件可能被各标签页再次广播,消费者必须保证幂等。若目标是减少连接数,需要额外实现 leader election,让主标签页持有 Socket,其他标签页只接收广播。


第二步:实现通用 Socket Store

下面是一份完整、可运行的 Store。它覆盖了:

  • SSR 环境保护;
  • 指数退避重连;
  • 页面隐藏、网络离线暂停;
  • 最多 100 条的发送队列;
  • 心跳;
  • 订阅、退订和尾沿防抖;
  • 消息过滤与归一化;
  • 跨标签页广播;
  • URL 敏感参数脱敏。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
// store/socket.js
import { create } from "zustand";
import { getBroadcastService } from "@/utils/broadcastService";

const CLOSED = 3;
const NORMAL_CLOSE = 1000;
const MAX_QUEUE_SIZE = 100;

const initialState = {
isConnected: false,
retryCount: 0,
visibility: "",
reason: "",
subscriptionsCount: 0,
readyState: CLOSED,
};

const inBrowser = () => typeof window !== "undefined";

const serialize = (message) =>
typeof message === "string" ? message : JSON.stringify(message);

export function sanitizeSocketUrl(url = "") {
try {
const parsed = new URL(url);
["jwt", "token", "deviceId", "device-id"].forEach((key) => {
if (parsed.searchParams.has(key)) {
parsed.searchParams.set(key, "[REDACTED]");
}
});
return parsed.toString();
} catch {
return url.replace(
/([?&](?:jwt|token|deviceId|device-id)=)[^&]*/gi,
"$1[REDACTED]",
);
}
}

const defaultNormalize = (event) => event;

export function createSocketStore({
ignoredMessages = ["PONG"],
normalizeMessage = defaultNormalize,
broadcast = false,
reportError,
} = {}) {
let socket = null;
let currentUrl = "";
let currentOptions = {};
let connectionVersion = 0;
let subscriptions = [];
let messageQueue = [];
let reconnectTimer = null;
let debounceTimers = new Map();
let broadcastService = null;
let removeBroadcastListener = null;

const getReadyState = () => {
if (!inBrowser()) return CLOSED;
return socket?.readyState ?? WebSocket.CLOSED;
};

const clearReconnectTimer = () => {
if (reconnectTimer) clearTimeout(reconnectTimer);
reconnectTimer = null;
};

const clearDebounceTimer = (id) => {
const timer = debounceTimers.get(id);
if (timer) clearTimeout(timer);
debounceTimers.delete(id);
};

const clearAllDebounceTimers = () => {
debounceTimers.forEach(clearTimeout);
debounceTimers.clear();
};

return create((set, get) => {
const updateStatus = (patch) => {
set({ ...patch, readyState: getReadyState() });
};

const updateSubscriptionsCount = () => {
set({ subscriptionsCount: subscriptions.length });
};

const distribute = (message) => {
[...subscriptions].forEach(({ id, handler, debounceDelay }) => {
clearDebounceTimer(id);

if (!debounceDelay) {
try {
handler(message);
} catch (error) {
console.error("WebSocket subscriber failed", error);
}
return;
}

const timer = setTimeout(() => {
try {
handler(message);
} catch (error) {
console.error("WebSocket subscriber failed", error);
} finally {
debounceTimers.delete(id);
}
}, debounceDelay);

debounceTimers.set(id, timer);
});
};

const ensureBroadcast = () => {
if (!broadcast || !inBrowser() || removeBroadcastListener) return;

broadcastService = getBroadcastService();
const listener = (data) => {
const normalized = normalizeMessage(
{ data, origin: window.location.origin },
data,
{ fromBroadcast: true },
);
if (normalized) distribute(normalized);
};

broadcastService.addListener(listener);
removeBroadcastListener = () => {
broadcastService?.removeListener(listener);
removeBroadcastListener = null;
};
};

const handleMessage = (event) => {
if (ignoredMessages.includes(event.data)) return;

let parsed;
try {
parsed = JSON.parse(event.data);
} catch {
parsed = undefined;
}

const normalized = normalizeMessage(event, parsed, {
fromBroadcast: false,
});
if (!normalized) return;

if (broadcast && broadcastService && parsed !== undefined) {
broadcastService.broadcast(parsed);
}
distribute(normalized);
};

const flushQueue = () => {
while (messageQueue.length && socket?.readyState === WebSocket.OPEN) {
const message = messageQueue.shift();
try {
socket.send(serialize(message));
} catch (error) {
console.error("WebSocket queued message failed", error);
}
}
};

const closeSocket = (
code = NORMAL_CLOSE,
reason = "Socket closed",
{ markStale = true } = {},
) => {
if (!socket) return;
const target = socket;
socket = null;
if (markStale) connectionVersion += 1;

try {
target.close(code, reason);
} catch (error) {
console.error("WebSocket close failed", error);
}

updateStatus({ isConnected: false, reason });
};

const connect = (url, options = {}) => {
if (!inBrowser() || !url || !navigator.onLine) return false;

currentUrl = url;
currentOptions = options;
ensureBroadcast();
clearReconnectTimer();

if (socket) closeSocket(NORMAL_CLOSE, "Reconnect");

const version = connectionVersion + 1;
const ws = new WebSocket(url);
connectionVersion = version;
socket = ws;
updateStatus({ isConnected: false, reason: "" });

ws.onopen = () => {
if (version !== connectionVersion) return;
updateStatus({
isConnected: true,
retryCount: 0,
visibility: document.visibilityState || "visible",
reason: "",
});
flushQueue();
options.onOpen?.();
};

ws.onmessage = (event) => {
if (version === connectionVersion) handleMessage(event);
};

ws.onclose = (event) => {
if (version !== connectionVersion) return;
if (socket === ws) socket = null;

const canReconnect =
event.code !== NORMAL_CLOSE && navigator.onLine && currentUrl;

if (!canReconnect) {
updateStatus({
isConnected: false,
reason: String(event.reason || event.code),
});
return;
}

const retryCount = get().retryCount + 1;
const delay = Math.min(1000 * 2 ** retryCount, 10000);
updateStatus({
isConnected: false,
retryCount,
reason: String(event.reason || event.code),
});

reconnectTimer = setTimeout(() => {
get().connect(currentUrl, currentOptions);
}, delay);
};

ws.onerror = (error) => {
if (version !== connectionVersion) return;
reportError?.(error, {
url: sanitizeSocketUrl(currentUrl),
readyState: ws.readyState,
retryCount: get().retryCount,
});
options.onError?.(error);
};

return true;
};

const suspendConnection = (reason = "Connection suspended") => {
clearReconnectTimer();
if (socket) closeSocket(NORMAL_CLOSE, reason);
else updateStatus({ isConnected: false, reason });
};

const disconnect = (reason = "Component unmounting") => {
suspendConnection(reason);
clearAllDebounceTimers();
removeBroadcastListener?.();
subscriptions = [];
messageQueue = [];
currentUrl = "";
currentOptions = {};
set({ ...initialState, reason });
};

return {
...initialState,
connect,
disconnect,
suspendConnection,

pauseForHiddenPage: () => {
updateStatus({ visibility: "hidden" });
if (
socket &&
[WebSocket.OPEN, WebSocket.CONNECTING].includes(socket.readyState)
) {
closeSocket(NORMAL_CLOSE, "Page hidden");
}
},

resumeVisiblePage: () => {
updateStatus({ visibility: "visible" });
if (
inBrowser() &&
navigator.onLine &&
currentUrl &&
(!socket || socket.readyState === WebSocket.CLOSED)
) {
get().connect(currentUrl, currentOptions);
}
},

subscribe: (handler, debounceDelay = 0) => {
const id = Symbol("socket-subscription");
subscriptions = [
...subscriptions,
{ id, handler, debounceDelay },
];
updateSubscriptionsCount();

return () => {
clearDebounceTimer(id);
subscriptions = subscriptions.filter((item) => item.id !== id);
updateSubscriptionsCount();
};
},

unsubscribe: (identifier) => {
const removed = subscriptions.filter(
(item) =>
item.id === identifier ||
item.handler === identifier ||
identifier === "*",
);
removed.forEach((item) => clearDebounceTimer(item.id));
subscriptions = subscriptions.filter((item) => !removed.includes(item));
updateSubscriptionsCount();
return removed.length;
},

sendMessage: (message) => {
if (socket?.readyState === WebSocket.OPEN) {
try {
socket.send(serialize(message));
set({ readyState: getReadyState() });
return true;
} catch (error) {
console.error("WebSocket send failed", error);
return false;
}
}

if (messageQueue.length >= MAX_QUEUE_SIZE) messageQueue.shift();
messageQueue.push(message);
set({ readyState: getReadyState() });
return false;
},

sendPing: () => {
if (
socket?.readyState === WebSocket.OPEN &&
navigator.onLine
) {
socket.send("PING");
return true;
}
return false;
},

forceReconnect: () => {
clearReconnectTimer();
updateStatus({ retryCount: 0 });
if (currentUrl) get().connect(currentUrl, currentOptions);
},

getSubscriptionsCount: () => subscriptions.length,
getReadyState,
};
});
}

const normalizeAssetMessage = (event, parsed) => {
if (parsed === undefined) return null;
return { data: parsed };
};

export const assetSocketStore = createSocketStore({
ignoredMessages: ["PONG", "CONNECTED"],
normalizeMessage: normalizeAssetMessage,
broadcast: true,
});

export const marketSocketStore = createSocketStore();
export const tradeSocketStore = createSocketStore();

const createSocketHook = (store) => () => ({
isConnected: store((state) => state.isConnected),
retryCount: store((state) => state.retryCount),
reason: store((state) => state.reason),
subscribe: store((state) => state.subscribe),
unsubscribe: store((state) => state.unsubscribe),
sendMessage: store((state) => state.sendMessage),
sendPing: store((state) => state.sendPing),
forceReconnect: store((state) => state.forceReconnect),
});

export const useAssetSocketApi = createSocketHook(assetSocketStore);
export const useMarketSocketApi = createSocketHook(marketSocketStore);
export const useTradeSocketApi = createSocketHook(tradeSocketStore);

这里有三个实现细节值得解释。

为什么不把 socket 放进 Zustand state?

WebSocket 是可变对象,也不需要渲染到 UI。放进 state 既不能获得不可变数据的收益,还会让组件有机会越过统一 API 直接修改连接。闭包更适合保存这类基础设施对象。

为什么使用 connectionVersion

快速切换 URL 时,旧连接的 oncloseonmessage 可能晚于新连接触发。版本号让旧连接事件自动失效,避免旧事件污染新页面。

订阅参数为什么叫 debounceDelay

收到新消息时先清除旧定时器,再等待一段时间执行,这是尾沿防抖,不是节流。若把它命名为 throttleDelay,调用者很容易误判高频行情的处理行为。


第三步:让 Bootstrap 接管 React 生命周期

Store 不应该依赖某个具体页面,但连接必须跟随 React 树挂载和卸载。WebSocketBootstrap 就是两者之间的桥梁。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
// components/Bootstrap/WebSocketBootstrap.jsx
import { useEffect, useRef } from "react";
import { marketSocketStore } from "@/store/socket";

export default function WebSocketBootstrap({
children,
url,
options = {},
store = marketSocketStore,
heartbeatInterval = 2000,
}) {
const optionsRef = useRef(options);
const isConnected = store((state) => state.isConnected);
const connect = store((state) => state.connect);
const disconnect = store((state) => state.disconnect);
const suspendConnection = store((state) => state.suspendConnection);
const pauseForHiddenPage = store((state) => state.pauseForHiddenPage);
const resumeVisiblePage = store((state) => state.resumeVisiblePage);
const sendPing = store((state) => state.sendPing);

useEffect(() => {
optionsRef.current = options;
}, [options]);

useEffect(() => {
if (!url) return undefined;
connect(url, optionsRef.current);

const handleVisibility = () => {
if (document.visibilityState === "visible") resumeVisiblePage();
else pauseForHiddenPage();
};
const handleOnline = () => resumeVisiblePage();
const handleOffline = () => suspendConnection("Network offline");

document.addEventListener("visibilitychange", handleVisibility);
window.addEventListener("focus", handleOnline);
window.addEventListener("online", handleOnline);
window.addEventListener("offline", handleOffline);

return () => {
document.removeEventListener("visibilitychange", handleVisibility);
window.removeEventListener("focus", handleOnline);
window.removeEventListener("online", handleOnline);
window.removeEventListener("offline", handleOffline);
disconnect("Component unmounting");
};
}, [
url,
connect,
disconnect,
suspendConnection,
pauseForHiddenPage,
resumeVisiblePage,
]);

useEffect(() => {
if (!isConnected || !heartbeatInterval) return undefined;
const timer = setInterval(sendPing, heartbeatInterval);
return () => clearInterval(timer);
}, [isConnected, heartbeatInterval, sendPing]);

return children;
}

页面隐藏时主动关闭连接,可以减少后台标签页的连接和消息处理成本;页面恢复、窗口重新聚焦或网络上线后再重连。是否这样做要结合业务判断:如果后台也必须持续接收消息,应保留连接,只降低 UI 更新频率。


第四步:在 _app.jsx 启动全局连接

市场连接不需要登录,可以始终挂载;资产连接只有登录后才需要。下面只展示与 Socket 相关的最小结构。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
// pages/_app.jsx
import { useEffect, useState } from "react";
import WebSocketBootstrap from "@/components/Bootstrap/WebSocketBootstrap";
import {
assetSocketStore,
marketSocketStore,
} from "@/store/socket";

function AppContent({ Component, pageProps, userSession }) {
const [fallbackId, setFallbackId] = useState("");

useEffect(() => {
setFallbackId(crypto.randomUUID());
}, []);

const marketId = pageProps.wsid || fallbackId;

const marketUrl = marketId
? `${process.env.NEXT_PUBLIC_MARKET_SOCKET}/${marketId}-market`
: null;
const assetUrl = userSession
? `${process.env.NEXT_PUBLIC_ASSET_SOCKET}?jwt=${encodeURIComponent(
userSession.accessToken,
)}&deviceId=${encodeURIComponent(userSession.deviceId)}`
: null;

return (
<WebSocketBootstrap
url={marketUrl}
store={marketSocketStore}
heartbeatInterval={2000}
>
<WebSocketBootstrap
url={assetUrl}
store={assetSocketStore}
heartbeatInterval={2000}
>
<Component {...pageProps} />
</WebSocketBootstrap>
</WebSocketBootstrap>
);
}

示例沿用了既有后端协议把凭证放在 URL 中,但从安全角度更推荐短期 Socket ticket 或握手后的鉴权消息。URL 可能出现在代理日志、浏览器工具和异常上下文中;如果协议暂时不能改,监控上报必须同时脱敏 token 与设备标识。

登录变为未登录时,assetUrl 从字符串变成 null。React 会先执行旧 effect 的 cleanup,因此旧连接仍会被 disconnect() 清理,然后新 effect 因没有 URL 而不连接。


第五步:在交易页面启动独立连接

交易页连接跟当前标的绑定,适合放在页面布局中,而不是全局 _app.jsx

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// pages/spot/[id].js
import WebSocketBootstrap from "@/components/Bootstrap/WebSocketBootstrap";
import { tradeSocketStore } from "@/store/socket";
import { TradeProvider } from "@/context/TradeContext";

export default function SpotPage({ id, info, wsid }) {
const tradeUrl = `${process.env.NEXT_PUBLIC_MARKET_SOCKET}/${wsid}`;

return (
<WebSocketBootstrap url={tradeUrl} store={tradeSocketStore}>
<TradeProvider id={id} type="spot" info={info}>
<SpotHeader />
<OrderBook />
<BuySell />
<AssetPanel />
</TradeProvider>
</WebSocketBootstrap>
);
}

这里形成了两种不同的状态边界:

  • tradeSocketStore 管理连接,是客户端模块级基础设施;
  • TradeProvider 内部创建页面级 vanilla Store,管理当前标的、最新价、深度和下单 UI 状态。

Socket 可以被同一页面的多个组件订阅,但交易状态不会泄漏到其他页面会话。


第六步:业务组件如何订阅消息

组件只需要关心三件事:连接是否就绪、发送什么订阅协议、如何处理消息。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import { useCallback, useEffect, useState } from "react";
import { useTradeSocketApi } from "@/store/socket";

export default function LatestPrice({ symbol }) {
const [price, setPrice] = useState("--");
const { isConnected, sendMessage, subscribe } = useTradeSocketApi();

const handleMessage = useCallback(
(event) => {
let message;
try {
message = JSON.parse(event.data);
} catch {
return;
}

if (message.type !== "ticker" || message.symbol !== symbol) return;
setPrice(message.lastPrice);
},
[symbol],
);

useEffect(() => {
if (!isConnected) return undefined;

// 每次重连成功都会重新发送,恢复服务端订阅。
sendMessage({ action: "subscribe", topic: "ticker", symbol });
const unsubscribeLocal = subscribe(handleMessage);

return () => {
unsubscribeLocal();
sendMessage({ action: "unsubscribe", topic: "ticker", symbol });
};
}, [isConnected, symbol, sendMessage, subscribe, handleMessage]);

return <strong>{price}</strong>;
}

subscribe() 返回清理函数,比再次用 handler 调 unsubscribe() 更稳妥,也不容易因为函数引用变化漏掉监听器。

还要区分两类订阅:

  1. 本地订阅:把组件 handler 注册到前端分发器;
  2. 服务端订阅:通过 sendMessage() 告诉服务器需要哪些 topic。

本地 unsubscribe() 不等于服务端退订。组件卸载时是否要发送退订消息,取决于服务端协议是增量订阅、覆盖订阅,还是连接关闭后自动清理。


列表场景:只订阅可视区域

市场列表可能有几百个标的,没有必要全部订阅。可以配合 IntersectionObserver,维护当前可见项并发送覆盖式订阅。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
import { useEffect, useMemo, useRef } from "react";
import { useMarketSocketApi } from "@/store/socket";

export function QuoteRow({ item, onVisibleChange }) {
const rowRef = useRef(null);

useEffect(() => {
const node = rowRef.current;
if (!node) return undefined;

const observer = new IntersectionObserver(([entry]) => {
onVisibleChange(item, entry.isIntersecting);
});
observer.observe(node);

return () => {
onVisibleChange(item, false);
observer.disconnect();
};
}, [item, onVisibleChange]);

return <tr ref={rowRef}>{/* 行情单元格 */}</tr>;
}

export function useVisibleQuoteSubscription() {
const visibleMap = useRef(new Map());
const timerRef = useRef(null);
const { isConnected, sendMessage } = useMarketSocketApi();

const flush = useMemo(
() => () => {
if (!isConnected) return;
sendMessage({
action: "replace-subscriptions",
symbols: [...visibleMap.current.keys()],
});
},
[isConnected, sendMessage],
);

const onVisibleChange = useMemo(
() => (item, visible) => {
if (visible) visibleMap.current.set(item.symbol, item);
else visibleMap.current.delete(item.symbol);

clearTimeout(timerRef.current);
timerRef.current = setTimeout(flush, 300);
},
[flush],
);

useEffect(() => {
if (isConnected) flush();
}, [isConnected, flush]);

useEffect(
() => () => {
clearTimeout(timerRef.current);
},
[],
);

return onVisibleChange;
}

覆盖式协议必须允许发送空数组,否则所有元素离开可视区域后,服务器仍可能保留上一次订阅。节流或防抖也不能简单地“时间不足就 return”,否则最后一次可见列表变化会永久丢失;应该保留尾沿执行。


重连时如何恢复订阅?

连接断开后,服务端通常会清掉这条连接上的 topic。前端必须在重连成功后恢复订阅。

这套方案没有在基础 Store 中维护业务 topic 注册表,因为不同服务端的协议可能完全不同:

  • 有的按单个 symbol 增量订阅;
  • 有的每次发送完整列表并覆盖旧列表;
  • 有的订阅消息还包含市场、产品类型和权限信息。

因此恢复动作放在业务组件里,通过 isConnected effect 重新发送。连接 Store 只保证:

  • 连接打开前调用 sendMessage() 会进入有界队列;
  • 连接异常关闭会自动重连;
  • isConnectedfalse 变为 true 时组件会重新执行订阅 effect。

如果项目协议已经统一,可以再把 topic registry 下沉到 Store,但不要让基础设施层猜测业务协议。


从 React Context 迁移到 Zustand 的过程

这次迁移没有一次性删除所有 Context,而是按职责拆分:

  1. 先提取通用连接能力,建立 createSocketStore()
  2. WebSocketBootstrap 替代原 WebSocket Provider;
  3. 为资产消息补充 { data: parsedObject } 兼容层;
  4. 保留资产跨标签页广播,并对监控 URL 脱敏;
  5. 将消费组件逐个从旧 Context Hook 替换到 useAssetSocketApi()
  6. 删除组件内部重复的心跳定时器;
  7. 保留资产业务 Context 和页面级 Trade Store,因为它们解决的是业务状态隔离,不是连接管理。

这种迁移方式的价值在于:调用方可以逐个替换,连接层和业务层不会同时重写。对于资产、交易、登录授权一类高风险链路,小步迁移比“大一统重构”更容易验证和回滚。


当前方案的优势与边界

已解决的问题

  • 三类 Socket 独立,页面生命周期互不干扰;
  • 心跳、断网、页面可见性、重连逻辑集中管理;
  • 发送队列有上限,不会在长期断网时无限增长;
  • 组件可以精确选择连接状态,不依赖庞大的 Context value;
  • 高频原始消息不进入 Zustand state,降低全局渲染压力;
  • 每个订阅都有明确清理路径;
  • 旧连接事件通过版本号失效;
  • 监控上报可以统一做 URL 脱敏。

仍需注意的问题

  1. BroadcastChannel 不是连接复用。 没有主标签页选举时,每个标签页仍会建立连接,还可能产生重复广播。
  2. 消息队列可能过期。 断网期间排队的订阅消息在恢复后未必仍符合当前页面,业务层应使用可替换的订阅快照或在切页时清理。
  3. 本地监听数不是服务端订阅数。 subscriptionsCount 只能反映前端 handler 数量。
  4. 控制消息过滤要基于协议。 不应因为响应里存在一个 code 字段就全部丢弃,业务错误和握手响应需要分别处理。
  5. 凭证不宜长期放在 URL。 至少要在日志、Sentry、埋点中脱敏 token 和设备标识。
  6. 隐藏即断开不是通用答案。 后台通知或必须连续采样的页面可能需要维持连接。
  7. 尾沿防抖不是节流。 高频行情需要固定频率采样时,应实现真正的 throttle,而不是反复推迟执行。

验证清单

Socket 改造不能只看“页面有数据”,至少需要覆盖下面这些场景:

  • 首次进入页面只建立预期数量的连接;
  • Socket 打开前发送的消息能在连接后按顺序发出;
  • 异常关闭后按指数退避重连,正常关闭不重连;
  • 重连后组件重新发送服务端订阅;
  • 页面隐藏、恢复可见、离线、上线行为符合产品预期;
  • URL 改变时旧连接事件不会污染新连接;
  • 组件卸载后本地订阅数归零,心跳和防抖定时器被清除;
  • 未登录时不建立资产连接,退出登录后旧资产连接关闭;
  • 正式环境与模拟环境不会串线;
  • 多标签页下没有重复业务副作用;
  • 监控、日志和埋点不包含 token、设备 ID 或用户隐私字段;
  • 无效 JSON、控制消息和业务 handler 异常不会导致整条连接停止分发。

总结

Socket + Zustand 的重点从来不是把 new WebSocket() 搬进 Store,而是建立清晰的边界:

  • 连接基础设施放在 Store 工厂闭包;
  • UI 需要的连接状态放在 Zustand;
  • 生命周期交给 Bootstrap;
  • 服务端订阅协议和业务消息解析留在业务层;
  • 页面会话状态使用独立 vanilla Store 隔离;
  • 所有定时器、监听器、订阅和连接都有对称清理。

当这几层职责稳定下来,新增一条 Socket 连接不再意味着复制一套 Provider,也不会迫使所有实时数据进入同一个全局状态树。Zustand 在这里真正提供的,是一个轻量、可选择订阅、又能脱离 React 组件存在的状态边界。