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

RTK Query中WebSocket断开重连后数据未更新的解决方案

问题描述

尝试用RTK Query建立WebSocket连接,当前代码会在连接断开时自动创建新连接,但重连后useGetMessagesQuery中的数据无法更新。

ChatApi代码:

let ws: WebSocket;
let interval: ReturnType<typeof setInterval>;
export const init = () => {
  ws = new WebSocket('ws://localhost:8080');
  ws.onopen = (message: Event) => {
    interval = setInterval(() => {
      ws.send(JSON.stringify({ type: 'keepAlive' }));
    }, 1000 * 15);
  };
  ws.onclose = (message: CloseEvent) => {
    interval && clearInterval(interval);
    setTimeout(() => {
      init();
    }, 1000 * 10);
  };
}

export const chatApi = createApi({
  baseQuery: fetchBaseQuery({ baseUrl: '/' }),
  endpoints: (build) => ({
    getMessages: build.query<any, any>({
      queryFn: () => ({ data: {} }),
      async onCacheEntryAdded(
        arg,
        { updateCachedData, cacheDataLoaded, cacheEntryRemoved }
      ) {
        try {
          await cacheDataLoaded
          const listener = (event: MessageEvent) => {
            const data = JSON.parse(event.data)
            updateCachedData((draft) => {
              draft.push(data)
            })
          }

          ws.addEventListener('message', listener)
        } catch {
        }
        await cacheEntryRemoved
        ws.close()
      },
    }),
  }),
})

export const { useGetMessagesQuery } = chatApi;

组件中使用消息数据:

const { data: message } = useGetMessagesQuery();

useEffect(() => {
  init();
}, []);

useEffect(() => {
  console.log(message);
}, [message]);

问题原因与修复方案

核心问题

重连后新创建的WebSocket实例没有重新绑定消息监听,原来的监听仍挂在旧的WebSocket对象上,导致新连接的消息无法触发缓存更新。此外,onCacheEntryAdded中关闭全局WebSocket的逻辑会干扰重连流程。

修复步骤

  1. 重构WebSocket管理逻辑,实现订阅机制
    不再用全局变量直接覆盖ws,而是维护消息监听列表,重连时自动触发所有订阅回调:

    let ws: WebSocket | null = null;
    let interval: ReturnType<typeof setInterval> | null = null;
    const messageListeners: ((event: MessageEvent) => void)[] = [];
    
    export const init = () => {
      if (ws?.readyState === WebSocket.OPEN) return;
      
      ws = new WebSocket('ws://localhost:8080');
      ws.onopen = () => {
        interval = setInterval(() => {
          ws?.send(JSON.stringify({ type: 'keepAlive' }));
        }, 15000);
      };
      ws.onmessage = (event) => {
        messageListeners.forEach(listener => listener(event));
      };
      ws.onclose = () => {
        interval && clearInterval(interval);
        setTimeout(() => init(), 10000);
      };
      ws.onerror = () => ws?.close();
    };
    
    export const subscribeToMessages = (listener: (event: MessageEvent) => void) => {
      messageListeners.push(listener);
      return () => {
        const index = messageListeners.indexOf(listener);
        if (index !== -1) messageListeners.splice(index, 1);
      };
    };
    
    export const chatApi = createApi({
      baseQuery: fetchBaseQuery({ baseUrl: '/' }),
      endpoints: (build) => ({
        getMessages: build.query<any[], void>({
          queryFn: () => ({ data: [] }),
          async onCacheEntryAdded(
            _,
            { updateCachedData, cacheDataLoaded, cacheEntryRemoved }
          ) {
            init();
            await cacheDataLoaded;
            
            const unsubscribe = subscribeToMessages((event) => {
              try {
                const data = JSON.parse(event.data);
                updateCachedData(draft => draft.push(data));
              } catch (e) {
                console.error('消息解析失败:', e);
              }
            });
    
            await cacheEntryRemoved;
            unsubscribe();
          },
        }),
      }),
    });
    
    export const { useGetMessagesQuery } = chatApi;
    
  2. 简化组件初始化逻辑
    无需在组件中手动调用init(),onCacheEntryAdded会自动初始化WebSocket:

    const { data: messages } = useGetMessagesQuery();
    
    useEffect(() => {
      console.log('当前消息:', messages);
    }, [messages]);
    
  3. 额外优化建议

    • 明确指定缓存数据类型(如any[]),避免类型错误
    • 重连后可主动请求历史消息,补充缓存数据
    • 增加消息解析的错误捕获,防止单个无效消息中断缓存更新

替代实现思路

可以借助RTK Query的cacheEntryRemoved生命周期管理WebSocket连接,或者结合lazyQuery动态控制连接时机,但核心逻辑都是确保重连时重新绑定消息监听,让新连接的消息能触发缓存更新。

内容的提问来源于stack exchange,提问作者user16223807

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 01:28:38