ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

FastGPT Workflow 节点响应持久化改造:Append-Only 存储与交互恢复 NodeResponse ID 设计解析

FastGPT Workflow 节点响应持久化改造:Append-Only 存储与交互恢复 NodeResponse ID 设计解析 FastGPT Workflow 节点响应持久化改造Append-Only 存储与交互恢复 NodeResponse ID 设计解析【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPTFastGPT 的 workflow 运行详情通过chat_item_responses集合平铺持久化每条 row 的data即一个节点响应nodeResponse。本文以仓库内权威设计文档node-response-append-only-interactive-id.md为主线结合 nodeResponseStorage.ts、nodeResponseSink.ts、mergeNode.ts 等实现源码系统讲解 append-only 数据模型、读取时的增量合并算法、交互恢复场景下 nodeResponse ID 的复用规则以及运行前preChatRound的职责边界。读完你将掌握 FastGPT 节点详情从“运行期可更新”迁移到“只追加 读取时折叠”的完整设计思路以及交互恢复如何避免展示节点重复。背景从“可更新存储”到“只追加存储”的演进workflow 运行详情通过chat_item_responses平铺保存。每条 row 的data是一个 nodeResponse包含三个关键身份字段data.id展示节点 ID标识一个节点响应实例data.parentId父展示节点 ID读取时用于还原childrenResponses树形结构chatItemDataId所属 AI chat item 的dataId即本轮响应消息 ID。早期设计依赖{ appId, chatId, chatItemDataId, data.id }唯一索引并在运行期先删除同data.id的旧 row 再写入新 row或用 replace 模式清空旧详情。该方案有两个核心痛点大表唯一索引成本高在承载海量节点详情的大表上维护复合唯一索引写入吞吐和锁竞争压力大违背运行期只追加的性能目标运行中频繁 delete/update 增加写放大且并行 retry 时“删除旧 rows 再写入”的时序很难保证一致性。因此当前权威方案将 nodeResponse 表调整为append-onlyworkflow 运行过程中只createrows不更新、不删除。重复展示节点不再依赖数据库去重而是通过读取时按(data.id, data.parentId)fold折叠合并来还原最终形态。该设计文档合并并替代了历史文档node-response-stream-persistence.md其中的data.idunique 索引、运行期 delete 后 create、replace/append 模式、parallel retry 删除旧 rows 等描述已过时以及旧版 append-only 讨论稿。核心结论速览设计文档沉淀的结论如下chat_item_responses运行期只追加 rows对话删除、应用删除、过期清理等外部清理流程可以批量删除。data.id不再是数据库唯一键只表示前端展示节点身份。同一个data.id且parentId相同的多条 rows 表示同一个展示节点的多次增量读取时合并成一个节点两条 row 都没有parentId时也视为同一个 parent。mergeSignId已废弃不再写入、不再读取、不再兼容旧合并语义。旧数据若依赖mergeSignId展示异常可接受迁移或回放另行处理。dispatchWorkFlow.responseChatItemId是必填运行参数dispatch 不生成兜底 ID也不查询MongoChatItem或MongoChatItemResponse判断是否重复。保存对话记录的新运行必须先走preChatRound由业务入口完成最终chatId/responseChatItemId解析、生成锁、AI dataId 冲突检查和 Human/AI placeholder 预创建。普通新运行中 Human 和 AI 使用同一个roundDataId responseChatItemId。Human/AI 同 dataId 是预期行为同一个obj下重复 dataId 才是不合法语义。本轮运行前只阻塞 AI dataId 冲突。数据模型与索引设计row 结构chat_item_responses的核心字段定义如下对应 chatItemResponseSchema.tstype ChatItemResponseSchema { teamId: ObjectId; appId: ObjectId; chatId: string; chatItemDataId: string; data: ChatHistoryItemResType; time: Date; };在真实 Schema 中appId字段注释说明了其历史物理字段名语义为sourceIdApp 场景才是真实 appIdsourceType来自ChatSourceTypeEnumtime默认为当前时间。保留的索引当前chat_item_responses保留两个索引ChatItemResponseSchema.index({ appId: 1, chatId: 1, chatItemDataId: 1, _id: 1 }); ChatItemResponseSchema.index({ teamId: 1, time: -1 });索引用途{ appId, chatId, chatItemDataId, _id }按 AI chat item 拉取完整 nodeResponse rows并按_id: 1保持写入顺序源码中复合索引包含_id避免详情读取时额外排序{ teamId, time: -1 }过期清理或团队维度清理。源码中还额外定义了一个{ sourceType, appId, chatId, chatItemDataId, _id }索引带 TODO 注释暂未全面检查操作故未加 sourceType 索引的完整方案说明数据访问正在向 sourceType 维度演进。明确不再创建的索引ChatItemResponseSchema.index( { appId: 1, chatId: 1, chatItemDataId: 1, data.id: 1 }, { unique: true } );chat_items当前保留普通索引ChatItemSchema.index({ appId: 1, chatId: 1, dataId: 1 }); ChatItemSchema.index({ appId: 1, chatId: 1, deleteTime: 1 }); ChatItemSchema.index({ appId: 1, chatId: 1, _id: -1 }); ChatItemSchema.index({ appId: 1, chatId: 1, obj: 1, _id: -1 });{ appId, chatId, dataId }不能改成 unique因为普通新运行中 Human 和 AI 会共享同一个dataId同一轮消息的 Human/AI 记录同 ID。如果后续 AI dataId 冲突检查需要优化可以补普通索引{ appId, chatId, dataId, obj }但不加 unique。写入路径WorkflowNodeResponseWriter写入封装在WorkflowNodeResponseWriternodeResponseStorage.ts其工作流程为一个 workflow 请求复用一个 writer子 workflow、loop、parallel、toolcall 等共享该 writer。record()接收本次要保存的 nodeResponses补齐id/parentId、裁剪 dataset quote、计算childResponseCount转成 flat rows对应createChatItemResponseRows。recordWithParent()只给没有parentId的 root child 补外层 parent已有parentId的响应保持内部层级避免破坏更细的层级结构。writer 通过 promise queuewriteQueue串行化并发record保证 Mongo_id顺序接近运行期写入顺序——因为子 workflow、parallel 分支可能并发调用同一个 writer串行化后才能保证详情展示顺序稳定。默认batchSize 5达到阈值或 close 时 flush。flush 只执行create(rowsWithTime, { ordered: true, session, ...writePrimary })不做任何 delete/update/replace。普通写入失败重试 3 次NODE_RESPONSE_WRITE_RETRY_TIMES 3仍失败则写 slim rows只保留节点身份、名称、类型、父子关系、运行时间和消耗统计等关键字段的瘦身版本slim 仍失败时丢弃本批详情 rows 并记录日志不阻断主 workflow。saveChat需要的引用citeCollectionIds、错误数和根节点积分由 writer 在运行期维护 summarysummaryContributionsMap按id parentId覆盖避免 retry/完成态重复累计详情 rows 写库失败不影响这些摘要。值得注意的是写入前不做 JSON/BSON 体积预估BSON 大小、不可序列化字段等问题统一交给 Mongo 写入校验失败后进入 retry/slim fallback避免正常路径额外 CPU 与临时内存开销。flush 后会立即释放 buffer降低运行期内存占用。另外数据集搜索节点datasetSearchNode的quoteList在入库前会被瘦身slimQuoteListForStorage只保留id/chunkIndex/datasetId/collectionId/sourceId/sourceName/score等引用关联、来源和分数元信息移除 q/a 完整文本——因为完整 quote 体积很大且详情展示只需要来源元信息瘦身可降低单条 row 过大导致 Mongo 写失败的概率。实时发布路径WorkflowNodeResponseSinkNodeResponse 的持久化和实时发布统一由请求级WorkflowNodeResponseSink协调nodeResponseSink.ts一个 workflow 请求只创建一个 sink内部复用同一个WorkflowNodeResponseWriter。root workflow、child workflow、Agent、ToolCall、LoopRun、ParallelRun共享该 sink。节点和 Agent adapter 只上交标准 nodeResponse不直接操作 writer也不直接发送flowNodeResponseSSE。sink 为缺少 parentId 的响应补调用方显式传入的 parentIdWorkflowNodeResponseInput.parentId调用 writer 规范化并写入再按请求可见性配置发布本次响应。writer 仍按batchSize批量物理写 Mongo“接收一个、返回一个”指每个逻辑 nodeResponse 都产生独立 SSE 事件不要求每条 response 单独执行 Mongo create。sink不负责RuntimeNodeResponseSummary、usage、计费、child count 或控制流判断这些仍由 WorkflowQueue/Agent collector 在各自运行作用域内计算避免跨作用域重复累计。同(id, parentId)的多条响应仍是 append-only 增量sink 不去重、不覆盖、不改变数值字段的增量语义。输出协议矩阵V2streamtrue, detailtrue可见 nodeResponse 逐条发送flowNodeResponse客户端按(id, parentId)拼树结束时不再发送完整 nodeResponse 数组。V1streamtrue, detailtrue运行期不发送单个 nodeResponse结束时一次性发送flowResponses。V1/V2streamfalse, detailtrue结束时在 JSONresponseData中一次性返回。V2 Share 流式完整 nodeResponse 逐条写库对外先按 public node/field 规则过滤再逐条发送为保持pushResult2Remote原有回调契约运行期间仍保留最终详情数组。Share 可见性分层处理Share 可见性必须分层处理不能只依赖一个字段过滤函数responseAllDatafalsesink 只发布 public node 类型和字段并保留客户端拼树需要的id/parentId。Share workflow 内部始终保留回答中的引用 IDwriter 始终接收完整 nodeResponse普通 API 保持retainDatasetCite原有语义。datasetquoteList入库时继续移除 q/a只保留引用关联、来源和分数等元信息。showCite控制公开 nodeResponse 是否包含quoteList关闭时 SSE 与非流式 JSON 都不返回quoteList但不改写 SSE 回答文本也不改变持久化数据。客户端没有 quoteList 时不展示引用之后重新开启配置并刷新 Share可以根据已保存的引用 ID 和 quote 元信息恢复展示。showRunningStatus控制flowNodeStatus/toolCall/toolParams/toolResponse等过程事件不直接禁止引用展示依赖的 publicflowNodeResponse。showSkillReferences继续由 Agent 输出链路控制并受showRunningStatus约束。showWholeResponse/showFullText/canDownloadSource继续由前端能力和详情/引用/文件接口鉴权sink 不替代这些权限检查。明确隐藏内部 workflow 的系统插件继续既不写入 child rows也不发布 child 事件只保留外层工具节点响应。pushResult2Remote不属于本次 SSE 改造范围继续使用运行期finalResponseData调用/shareAuth/finish不增加延迟读库或回调协议变化。运行期明确删除的行为不按data.iddelete 旧 rows不做updateOne upsert不做 replace 模式不在持久化 buffer 中按data.id去重不依赖data.idunique 索引不为 retry 预生成 row_id做幂等极低概率重复 create 产生的冗余 rows 由读取 fold 吸收。persistToDb false 的场景persistToDb false的 writer 不写 Mongo只保留 summary 和可选内存详情retainInMemory适用于 debug、eval、临时运行等不保存历史的入口。这类入口仍必须给 dispatch 传随机responseChatItemId只是该 ID 不参与数据库查重。读取与合并按 (data.id, parentId) 折叠增量读取时先按 chat item 拉 rowsMongoChatItemResponse.find( { appId, chatId, chatItemDataId }, { data: 1 } ).sort({ _id: 1 });然后composeNodeResponseDetail()调用mergeNodeResponseDataByIdAndParent()做 fold实现见 mergeNode.ts规则如下只处理存在data.id的 rows无 id 的 row 会被丢弃无法参与合并。合并 identity 是(data.id, data.parentId)parentId不存在时归一为同一个空值getNodeResponseIdentityKey用\u0000分隔 id 与 parentId。同 identity 的多条 rows 合并为一个展示节点。数值字段按增量累加包括runningTime保留两位小数、totalPoints、childResponseCount、tokens含 input/output/toolCall/embedding/reRank/extension等。llmRequestIds去重合并。compressTextAgent、deepSearchResult这类结构化用量字段按现有规则累加。普通标量字段以后到的 incoming 为准。childrenResponses递归按同一规则合并。child row 早于 parent row 到达时先作为临时 rootparent 到达后回收挂到childrenResponses对应appendNodeResponseByParent的 orphan 回收逻辑。批量读取时使用mergeNodeResponseListByParent一次性挂树算法先按id parentId合并同层增量再按 parentId 挂到childrenResponses避免每条 row 递归扫描已构建的整棵树在 loop/parallel 产生大量 rows 时把复杂度从接近 O(n²) 降到以线性扫描为主。历史兼容边界新数据统一使用childrenResponsespluginDetail/toolDetail/loopDetail/parallelDetail/loopRunDetail只作为历史 detail 字段读取和递归统计来源getChildrenResponses会把这些旧字段与childrenResponses一并收集不再作为新链路的通用写入结构chat_items.responseData已废弃。读取时如果独立表没有 rows才回退旧内联详情getChatItemResponseData的 fallback 逻辑避免历史数据被空结果覆盖childTotalPoints不再对外保留mergeNodeResponseDataByIdAndParent最后会stripChildTotalPoints子节点积分展示由客户端基于childrenResponses现场计算。NodeResponse ID 语义与交互恢复普通节点随机 ID普通节点首次运行时生成随机data.idgetNanoid()。这类 ID 不需要可预测也不需要数据库唯一约束。交互恢复复用暂停前 ID交互恢复时需要复用暂停前记录的 nodeResponse ID避免同一个展示节点在恢复后拆成两个节点const nodeResponseId lastInteractive?.nodeResponseId lastInteractive.entryNodeIds?.includes(node.nodeId) ? lastInteractive.nodeResponseId : getNanoid();WorkflowInteractiveResponseType增加通用字段定义于 interactive/type.tsnodeResponseId?: string;该字段与entryNodeIds平级表示触发本次暂停的当前 workflow 节点对应的 nodeResponsedata.id。同一时间只允许一个暂停模式因此一个字符串即可表示当前恢复入口。嵌套交互嵌套交互继续沿用childrenResponse每一层 interactive 都可以携带自己的nodeResponseId。例如 ToolCall 包装的子 workflow 暂停时{ type: toolChildrenInteractive, entryNodeIds: [toolCallNodeId], nodeResponseId: tool-call-node-response-id, params: { childrenResponse: { type: userInput, entryNodeIds: [formNodeId], nodeResponseId: form-node-response-id }, toolParams: { toolCallId: call_xxx } } }恢复时ToolCall 节点复用toolChildrenInteractive.nodeResponseId子 workflow 复用childrenResponse.nodeResponseId新增 rows 继续写到同一条 AI chat item 的chatItemDataId下读取时父 ToolCall 和子节点都按(data.id, parentId)合并页面只展示一个 ToolCall 节点用量和运行时间按增量累加。LoopRun 恢复iteration wrapper 的 ID 派生LoopRun 的 iteration wrapper 是虚拟展示节点ID 由 loopRun 父 nodeResponse ID 派生id ${loopRunNodeResponseId}:iter:${iteration};这样同一个 loop 节点在不同父作用域下运行不会因为node.nodeId iteration冲突。交互恢复时只要 loopRun 父节点复用interactive.nodeResponseId同一轮 iteration wrapper 也会自然复用同一个data.id。LoopRun 暂停时会写一次当前 iteration wrapper作为暂停前 child nodeResponses 的 parent并把pendingIterationSummary存到 interactive params。恢复后同一个 wrapper ID 再写本次 resume 的增量统计。由于读取会累加数值字段恢复后的 wrapper 必须只写本次 resume 片段的增量值不能写暂停前后合并后的累计值——这是防止数值双算的关键约束文档在“后续关注”中明确要求 LoopRun、ToolCall 等恢复场景必须持续保证写入的是本次运行片段增量。运行前 preChatRound业务入口的职责边界保存历史的新运行进入 workflow 前只调用preChatRound实现见 prepare.ts。它负责解析最终chatId。空chatId自动生成随机 chatIdgetNanoid(24)NO_RECORD_HISTORIES即NO_RECORD_CHAT_ID NO_RECORD_HISTORIES表示不保存历史。解析最终responseChatItemId。请求未传时生成随机 ID。判断是否持久化 chat items 和 nodeResponse rows。持久化运行占用MongoChat.chatGenerateStatus generating。普通新运行检查 AIdataId冲突。普通新运行严格创建本轮 Human AI placeholder。交互继续复用上一条 AI 的dataId不创建新的 Human/AI placeholder。失败时如果已经占用生成状态立刻置为error。返回值type PreChatRoundResult { chatId: string; responseChatItemId: string; shouldPersistChatRound: boolean; shouldFinalizePreparedRound: boolean; };持久化判断统一为const finalChatId chatId NO_RECORD_CHAT_ID ? chatId : chatId || getNanoid(24); const shouldPersistChatRound finalChatId ! NO_RECORD_CHAT_ID;入口后续必须使用preparedRound.chatId和preparedRound.responseChatItemId不能继续使用请求里的原始值。nodeResponseWriteConfig.persistToDb应等于preparedRound.shouldPersistChatRound。普通新运行顺序解析最终chatId/responseChatItemIdNO_RECORD_HISTORIES直接返回不持久化结果不占用生成锁调用tryStartGenerateChat占用生成锁已有 generating 时抛ChatErrEnum.chatIsGenerating校验已有 AI chat item 中不存在同responseChatItemId严格 create 本轮 Human AI placeholder二者使用同一个dataId responseChatItemIdprepareChatRound使用严格 create 而非 upsert检查或创建失败时写生成状态error并抛错创建成功后才进入 workflow。AI dataId 冲突检查口径MongoChatItem.findOne( { appId, chatId, obj: ChatRoleEnum.AI, dataId: responseChatItemId }, dataId );只检查 AI 的原因Human/AI 同dataId是新运行的正常结构本轮 nodeResponse rows 归属于 AIchatItemDataIdHuman 历史重复不影响 nodeResponse append-only 的安全性可以离线审计不作为运行前阻塞条件。preChatRound保持在业务入口不下沉到dispatchWorkFlow。dispatch 被 debug、skill debug、MCP、outLink、定时触发等入口复用不应该感知source/sourceName/shareId/outLinkUid/userContent等 chat 保存字段。此外源码中stripUserContentFileUrls会清理用户消息里的文件临时 URL只保留 file key 参与持久化避免历史记录保存过期访问地址。删除与清理运行期 writer 不删除 nodeResponse rows。外部删除规则删除整条对话或批量日志时可以按chatId删除MongoChatItemResponse局部消息删除继续保持MongoChatItem软删除语义新数据 Human/AI 同dataId删除一轮消息时前端可以继续收集 Human 和 AI 的 dataId但发请求前应去重删除接口支持 bodycontentIdsbody 优先body 为空时兼容 querycontentId。OpenAPI 默认声明 body。客户端约束普通新运行一轮只生成一个roundDataIdHuman/AI 共用该值交互继续复用上一条 AIdataId不是新一轮 Human/AIReact list key不能只用dataId因为 Human/AI 可能相同应包含obj或_id/id前端按dataId更新 AI 记录时需要带 AI 语义避免命中同 ID HumannodeResponse SSE 合并和详情弹窗都应使用(id, parentId)合并语义不再依赖mergeSignId。测试要求与回归保障仓库为本次改造配套了完整的测试覆盖核心测试文件包括 nodeResponseStorage.test.ts、nodeResponseSink.test.ts、index.persistence.test.ts 与 mergeNode.test.ts。preChatRound相关普通新运行成功创建MongoChat、Human、AI placeholderHuman/AI 同dataId responseChatItemId初始responseChatItemId命中已有 AI直接抛错不进入 workflow不随机兜底初始responseChatItemId只命中 Human不按重复 ID 报错生成锁冲突抛ChatErrEnum.chatIsGenerating不创建 placeholderplaceholder 创建失败或重复校验失败生成状态置为error空chatId自动生成随机 chatId 并保存记录NO_RECORD_HISTORIES不写 chat、不写 chat item、不占用生成锁仍返回 dispatch 可用的随机responseChatItemId非 query 交互继续复用上一条 AIdataId不创建新 placeholder找不到上一条 AI 时抛错interactive query按新一轮创建 Human/AI placeholderfinalizeChatRound能在 Human/AI 同 dataId 时按obj更新两条记录。nodeResponse append-only 相关writer 写入只调用 create不执行运行期 delete/update/replacebuffer 中同data.id多条 rows 全部写入不预去重retry 保留 3 次普通重试和 slim fallback不依赖预生成_id读取按_id顺序 fold同(data.id, parentId)合并为一个展示节点相同data.id但不同parentId不合并parentId都不存在时视为同 parent 并合并数值字段按增量累加标量以后到为准llmRequestIds去重child 先于 parent 到达时最终能挂回 parentmergeSignId不参与合并。交互恢复相关暂停时interactive.nodeResponseId写入当前节点data.id交互继续时恢复入口复用interactive.nodeResponseId页面只展示一个节点ToolCall 子 workflow 暂停后继续父 ToolCall 和子 workflow 分别复用对应层级nodeResponseIdLoopRun 暂停后继续父 loopRun 复用interactive.nodeResponseIditeration wrapper 使用${loopRunNodeResponseId}:iter:${iteration}恢复后只写本次片段增量避免数值双算。索引回归ChatItemResponseSchema不声明{ appId, chatId, chatItemDataId, data.id }unique 索引保留{ appId, chatId, chatItemDataId, _id }读取索引ChatItemSchema.index({ appId, chatId, dataId })保持普通索引不改 unique如新增{ appId, chatId, dataId, obj }也只能是普通索引。后续关注与实施边界设计文档同时记录了需要持续关注的运维与演进事项append-only 会增加 rows 数量需要依赖对话删除、应用删除和过期清理控制表规模历史mergeSignId数据不迁移异常展示风险已接受如果线上 AI dataId 冲突检查成为热点再评估普通索引{ appId, chatId, dataId, obj }LoopRun、ToolCall 等恢复场景必须持续保证写入的是本次运行片段增量而不是累计值。本次实施 TODO 清单均已勾选完成表明改造范围是新增请求级WorkflowNodeResponseSink统一 writer 与 V2 SSE 发布WorkflowQueue 的 root/child runtime 通过 sink 逐条发布 nodeResponseAgent collector 移除 writer 依赖改为上交 sinkLoopRun/ParallelRun 虚拟任务节点改走 sink系统插件内部 workflow 使用禁用 sink 的作用域保持隐藏语义V1、非流式 JSON、Share public 过滤和引用/文件权限保持兼容最终全量测试中全仓并发出现 4 个 20 秒超时相关文件单独复跑全部通过。总结FastGPT 的 nodeResponse 持久化从“唯一索引 运行期更新”演进为“append-only 写入 读取时按 (data.id, parentId) 折叠”本质上是把“写时去重”的复杂度转移到了“读时合并”从而换取运行期稳定、低成本的只追加写入。配合请求级WorkflowNodeResponseSink统一持久化与 SSE 发布、preChatRound在业务入口完成 chat 语义校验与 placeholder 预创建、交互恢复时复用nodeResponseId保证展示节点唯一这套设计同时解决了大表写入性能、并行运行写入顺序、Share 可见性分层以及暂停/恢复场景的展示一致性问题。对于需要深入理解 FastGPT workflow 运行链路或设计类似对话式 AI 工作流引擎持久化方案的开发者建议进一步阅读 nodeResponseStorage.ts、nodeResponseSink.ts、mergeNode.ts 以及 prepare.ts 中对应的测试用例。【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表