You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 10:45:35