RTK Query最佳实践:通过单个WebSocket连接更新所有injectedEndpoints缓存
单WebSocket连接结合RTK Query的实现与缓存更新
我刚接触RTK Query,需要为整个应用仅使用一个WebSocket连接,参照GitHub示例写了下面的代码。现在要实现两个核心需求:
- 通过这个WebSocket订阅发送payload
- 收到WebSocket消息时,更新其他injectedEndpoints的缓存
import { ApiSlice } from 'api'; import { instrumentsAdapter } from './marketSlice'; const socket = new WebSocket(process.env.REACT_APP_SOCKET_BASE_URL); const socketConnected = new Promise((resolve, reject) => { // Connection opened try { socket.addEventListener('open', (event) => { resolve(event); }); } catch (err) { console.log('err', err); reject(err); } }); export const socketApi = ApiSlice.injectEndpoints({ endpoints: (builder) => ({ socketChannel: builder.mutation({ async queryFn(arg) { await socketConnected; const { type, topic } = arg; const sendPayload = { type, path: topic }; socket.send(JSON.stringify(sendPayload)); return { data: { messages: [] } }; }, async onCacheEntryAdded(arg, { cacheDataLoaded, cacheEntryRemoved }) { console.log('arg', arg); await cacheDataLoaded; // Listen for messages socket.onmessage = (res) => { const message = JSON.parse(res.data); try { // ApiSlice.util.updateQueryData('getInstrumentByRefId', arg, (draft) => { // console.log('arg', arg); // draft = { ...message.value, baseVolume: 3 }; // }); } catch (err) { console.log('err', err); } }; await cacheEntryRemoved; socket.close(); } }) }) }); export const { useSocketChannelMutation } = socketApi;
问题修复与优化实现
1. 解决连接复用与生命周期问题
当前代码每次调用mutation后都会关闭连接,无法实现单连接复用。可以通过订阅计数来管理连接生命周期:
import { ApiSlice } from 'api'; import { instrumentsAdapter } from './marketSlice'; const socket = new WebSocket(process.env.REACT_APP_SOCKET_BASE_URL); let subscriptionCount = 0; // 用计数控制连接关闭时机 const socketConnected = new Promise((resolve, reject) => { socket.addEventListener('open', (event) => { resolve(event); }); // 正确监听连接错误 socket.addEventListener('error', (err) => { console.log('WebSocket连接错误', err); reject(err); }); }); // 全局消息监听,避免重复覆盖onmessage socket.onmessage = (res) => { const message = JSON.parse(res.data); // 根据消息内容匹配要更新的缓存 if (message.value?.refId) { try { // 更新getInstrumentByRefId的缓存 ApiSlice.util.updateQueryData( 'getInstrumentByRefId', message.value.refId, // 匹配目标查询的参数 (draft) => { // Immer直接修改draft属性即可,无需重新赋值 Object.assign(draft, message.value); draft.baseVolume = 3; // 自定义修改逻辑 } ); } catch (err) { console.log('缓存更新失败', err); } } };
2. 修正Mutation逻辑
调整mutation的订阅管理逻辑,确保只有所有订阅取消后才关闭连接:
export const socketApi = ApiSlice.injectEndpoints({ endpoints: (builder) => ({ socketChannel: builder.mutation({ async queryFn(arg) { await socketConnected; const { type, topic } = arg; const sendPayload = { type, path: topic }; socket.send(JSON.stringify(sendPayload)); subscriptionCount++; // 订阅计数加1 return { data: { success: true } }; }, async onCacheEntryAdded(arg, { cacheEntryRemoved }) { await cacheEntryRemoved; subscriptionCount--; // 无活跃订阅时才关闭连接 if (subscriptionCount === 0) { socket.close(); } } }) }) }); export const { useSocketChannelMutation } = socketApi;
3. 关键注意点
- 不要覆盖onmessage:全局设置一次消息监听即可,每次调用mutation覆盖onmessage会导致只有最后一次订阅能接收消息
- Immer更新规则:
updateQueryData中的draft是Immer代理对象,直接修改属性即可,不要用draft = {...}重新赋值 - 连接生命周期:通过订阅计数确保单连接在所有组件取消订阅后才关闭,实现真正的复用
- 错误处理:
addEventListener不会抛出同步错误,应该监听socket的error事件处理连接异常
内容的提问来源于stack exchange,提问作者Amir Rezvani
相关产品推荐
相关产品推荐

