TanStack Query 与 SSE 流式数据集成

AI 对话场景下,TanStack Query 管理已持久化数据,SSE 流式运行时管理临时数据,两者通过 cancelQueries + setQueryData + invalidate 协调。

#type / synthesis #status / growing #tech / dev / frontend #resource / react

[!info] related notes

TanStack Query 与 SSE 流式数据集成

范围

AI 对话页面同时存在两种数据流:已持久化的服务端数据(适合 TanStack Query)和正在流式生成的临时数据(适合组件 runtime)。本文解决两者如何共存、何时同步、如何防止冲突。

为什么要放在一起理解

单独用 TanStack Query 没问题,单独用 SSE 也没问题。但两者共存时会出现一个关键冲突:

1. session query 发起请求,拿的是旧数据
2. SSE 先返回新 extractedInfo,写入 query cache
3. 第 1 步的旧请求后返回
4. 旧 response 覆盖 cache
5. 右侧 InfoPanel 数据倒退

不理解这个冲突,就会在”流式更新正常”和”偶尔数据倒退”之间反复排查。

两种消息的边界

已持久化消息 → 来自后端 API → 放 query cache
正在流式生成的消息 → 前端 runtime 临时状态 → 放 AssistantChatPanel 内部
  • query.data.messages = 已保存的历史消息
  • AssistantChatPanel runtime = 当前正在 streaming 的临时消息
  • 消息真正持久化后,再更新 query cache

关键:不要每个 token 都写 query cache。 流式 token 由聊天 runtime 内部维护,query cache 只同步最终持久化结果。

防止 SSE 与 query cache 冲突的四步策略

第一步:session query 设置较长 staleTime

const consultationQuery = useQuery({
  queryKey: consultationKeys.session(routeConversationId),
  queryFn: () => consultationApi.getConsultation(routeConversationId!),
  enabled: !!routeConversationId,
  staleTime: 60_000,          // 60 秒内不自动重新请求
  refetchOnWindowFocus: false, // 切换窗口也不重新请求
});

减少自动 refetch 的频率,但不能完全杜绝。

第二步:SSE 活跃期间禁用 session 自动 refetch

const [isStreaming, setIsStreaming] = useState(false);

const consultationQuery = useConsultationSessionQuery(routeConversationId, {
  enabled: !isStreaming,
});

注意:不要让 enabled: !isStreaming 影响页面初次加载。isStreaming 只在 SSE 过程中为 true。

第三步:SSE 写 cache 前取消在途请求

const handleExtractedInfoUpdate = useCallback(
  async (info: ExtractedInfo[]) => {
    if (!routeConversationId) return;

    // 关键:先取消当前 session query 的在途请求
    await queryClient.cancelQueries({
      queryKey: consultationKeys.session(routeConversationId),
    });

    // 再写入 SSE 增量结果
    queryClient.setQueryData(
      consultationKeys.session(routeConversationId),
      (old) =>
        old
          ? { ...old, extracted_info: info }
          : {
              extracted_info: info,
              phase: 'collecting',
              diagnosis: null,
              treatment_plan: null,
              pending_interactions: [],
            },
    );
  },
  [routeConversationId, queryClient],
);

cancelQueries 确保旧请求不会在 SSE 之后返回并覆盖新数据。

第四步:SSE 结束后做最终同步

const handleStreamDone = useCallback(() => {
  setIsStreaming(false);

  // 流结束后,以后端最终持久化结果为准
  queryClient.invalidateQueries({
    queryKey: consultationKeys.session(routeConversationId),
  });
}, [routeConversationId, queryClient]);

依赖路径

SSE 活跃中
  → 前端 SSE 增量结果为准
  → cancelQueries 防止旧请求覆盖
  → setQueryData 写入增量

SSE 结束
  → invalidateQueries
  → 后端最终持久化结果为准

这个边界非常清晰:流中以前端为准,流结束后以后端为准。

消息持久化同步

当 SSE 返回 messagePersisted 事件时,更新 query cache 里的消息 ID:

const handleMessagePersisted = useCallback(
  (clientMessageId: string, messageId: string) => {
    queryClient.setQueryData(
      consultationKeys.conversation(routeConversationId),
      (old) => {
        if (!old) return old;
        return {
          ...old,
          messages: old.messages.map((m) =>
            m.id === clientMessageId ? { ...m, id: messageId } : m,
          ),
        };
      },
    );
  },
  [queryClient, routeConversationId],
);

SSE 事件到缓存的映射

SSE 事件写入位置操作说明
streaming tokenAssistantChatPanel runtime内部 state不写 query cache
messagePersistedconversation query cachesetQueryData 更新消息 IDclientMessageId → serverId
titleGeneratedlist + detail cachesetQueryData 两处同时更新列表和详情
extractedInfoUpdatesession query cachecancelQueries + setQueryData先取消旧请求再写入
phaseChangesession query cachesetQueryData直接更新 phase
stream finishedsession query cacheinvalidateQueries最终和后端对齐

对比与易混淆点

场景正确做法错误做法
流式 token 更新runtime 内部维护每个 token 写 query cache
extractedInfo 增量cancelQueries + setQueryData直接 setQueryData(可能被旧请求覆盖)
流结束invalidateQueries不做最终同步(可能和后端不一致)
消息持久化setQueryData 更新 IDsetMessages(主来源应该是 query cache)
title 更新同时 set list + detail cache只更新其中一处(两处缓存不一致)
创建于 2026/7/2 更新于 2026/7/15