TanStack Query 与 SSE 流式数据集成
AI 对话场景下,TanStack Query 管理已持久化数据,SSE 流式运行时管理临时数据,两者通过 cancelQueries + setQueryData + invalidate 协调。
#type / synthesis
#status / growing
#tech / dev / frontend
#resource / react
[!info] related notes
- 所属 MOC: TanStack Query 知识地图
- 基础: TanStack Query 服务端状态
- 基础: 前端 SSE 消费设计
- 架构: 四层状态架构
- 缓存: TanStack Query 缓存失效模式
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 token | AssistantChatPanel runtime | 内部 state | 不写 query cache |
messagePersisted | conversation query cache | setQueryData 更新消息 ID | clientMessageId → serverId |
titleGenerated | list + detail cache | setQueryData 两处 | 同时更新列表和详情 |
extractedInfoUpdate | session query cache | cancelQueries + setQueryData | 先取消旧请求再写入 |
phaseChange | session query cache | setQueryData | 直接更新 phase |
stream finished | session query cache | invalidateQueries | 最终和后端对齐 |
对比与易混淆点
| 场景 | 正确做法 | 错误做法 |
|---|---|---|
| 流式 token 更新 | runtime 内部维护 | 每个 token 写 query cache |
| extractedInfo 增量 | cancelQueries + setQueryData | 直接 setQueryData(可能被旧请求覆盖) |
| 流结束 | invalidateQueries | 不做最终同步(可能和后端不一致) |
| 消息持久化 | setQueryData 更新 ID | setMessages(主来源应该是 query cache) |
| title 更新 | 同时 set list + detail cache | 只更新其中一处(两处缓存不一致) |