
RivetKit Workflow Engine 持久化内幕KV 批处理 flush 机制与 BARE Schema 同步指南【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors本篇指南聚焦 RivetKit 的rivetkit/workflow-engine包在 Actor KV 之上的持久化实现storage.flush(...)如何把工作流历史写入 Actor KV 的批处理上限128 条 / 976 KiB、为何要等所有写入与删除都成功后才清除 dirty 标记以及持久化 BARE Schema 如何与 RivetKit 主包保持镜像同步并重新构建。读完你会掌握 workflow engine 存储层的分块chunking语义、原子批处理的取舍、删除操作的并发策略以及一套可执行的 schema 版本化与同步工作流可直接用于排查写入超限问题和参与该引擎的持久化格式演进。一、背景为什么 flush 是 workflow engine 的命脉RivetKit Workflow Engine 是一个面向 TypeScript 的持久化durable执行引擎其核心模型是把工作流写成普通 async 函数但每一次操作都会记录进持久化的历史history中进程崩溃、Actor 被驱逐eviction或休眠后重新从历史中重放replay而不是重新执行。这一机制的完整架构说明见 architecture.md。从架构上讲持久化层由三层组成WorkflowContext / 运行循环调用ctx.step()、ctx.loop()、ctx.sleep()等 API产生条目Entry。Storage内存中的工作流状态表示包含 nameRegistry、history、entryMetadata、workflow state/output/error 等storage.ts 中的Storage结构与createStorage()。EngineDriver持久化后端接口KV 调度每个工作流实例独占一个隔离的 KV 命名空间driver.ts 的EngineDriver接口及 architecture.md 的 Isolation Model。Storage 与 Driver 之间的唯一桥梁就是flush()它把内存中变脏的数据成批写回 Driver 的 KV。本文的核心主题——KV 批处理上限与 dirty 标记清除时机——正是 CLAUDE.md 中第一条维护指引的内容storage.flush(...)chunks driver batch writes to actor KV limits (128 entries / 976 KiB payload) and clears dirty markers only after all write/delete operations succeed.下面逐层拆解这句话背后的实现。二、flush 的批处理语义128 条 / 976 KiB 的硬约束2.1 两个常量的来源在 storage.ts 中定义了三个关键常量export const MAX_KV_BATCH_ENTRIES 128; export const MAX_KV_BATCH_PAYLOAD_BYTES 976 * 1024; /** Max delete ops (one transaction/permit each) run at once, under the 128-permit cap. */ export const MAX_CONCURRENT_DELETES 64;MAX_KV_BATCH_ENTRIES 128单次 KV batch 最多包含 128 条写入。这是 Actor KV如 Cloudflare Durable Objects 风格的单事务写限制的条目数上限。MAX_KV_BATCH_PAYLOAD_BYTES 976 * 1024976 KiB单次 batch 的键值字节总和上限。注意它刻意留出余量——976 KiB而不是整 1 MiB是为了给 KV 事务自身的元数据开销如每个 key 的哈希、事务记录头留出空间避免恰好压线导致底层存储拒绝。MAX_CONCURRENT_DELETES 64删除操作的并发批次上限。由于一次删除事务同样占用一个permit与 128 条目共享同一配额模型引擎将删除并发控制在 64始终低于 128 的配额上限。2.2 flush 的完整执行流程flush()storage.ts的执行步骤可以归纳为收集 name registry 新增项只从flushedNameCount之后的索引开始写只 flush 上次之后新增的名字每个名字对应一个buildNameKey(i)写入。收集 dirty 条目遍历storage.history.entries仅对entry.dirty true的条目序列化后放入writes。收集 dirty 元数据遍历storage.entryMetadata仅对metadata.dirty的条目写入buildEntryMetadataKey(id)。收集工作流级状态state、output、error各自与flushed*快照比较发生变化才写入error 按 name/message 比较因为对象不是引用相等。写入如果driver.atomicBatch为 true整个writes作为一次 batch 提交否则用splitBatchWrites切成不超过 128 条 / 976 KiB 的多个 chunk 依次提交。执行待处理删除如果调用方传入了pendingDeletions来自循环历史裁剪 collectLoopPruning在 batch 写入之后执行runDeletes。只有全部成功后才清除 dirty写成功后才把 dirty 条目/元数据的dirty置回false并推进flushedNameCount、flushedState、flushedOutput、flushedError。通知观察者若期间发生了任何历史变更historyUpdated回调onHistoryUpdated()。2.3 分块算法 splitBatchWritessplitBatchWrites 是chunks driver batch writes的直接实现以MAX_KV_BATCH_ENTRIES128 条为第一上限以MAX_KV_BATCH_PAYLOAD_BYTES976 KiB为第二上限按key.byteLength value.byteLength累计任何单条写入如果本身超过 976 KiB直接抛错KV batch write is N bytes, exceeding the 976 KiB byte limit而不是静默拆分——因为单键超出上限无论怎么切都无法写入尽早失败比写入一半更安全输出为多个 chunk每个 chunk 内既不超条数也不超字节数。2.4 原子批处理atomicBatch 开关driver.ts 中EngineDriver声明了一个可选标志/** * Requires each logical storage flush to reach batch as one indivisible * unit. The driver must reject an oversized unit instead of splitting it. */ readonly atomicBatch?: boolean;atomicBatch的含义是该驱动后端要求一次逻辑 flush 必须作为一个不可分割的整体到达batch()驱动必须拒绝超限单元而不是拆分。此时 workflow engine 不再自行分块因为拆分会破坏后端事务的原子性而是把全部写入一次性交给驱动由后端自己决定如何处理超限。反之当驱动没有声明原子批处理时引擎用splitBatchWrites自行分块。这也是chunks driver batch writes这句指引的另一半含义分块与否取决于驱动是否支持原子批处理。对应测试 storage.test.ts 精确验证了这两种路径does not split a driver-declared atomic flush构造 129 条 name超过 128在AtomicRecordingDriver声明 atomicBatch下 flush 后driver.batches只有 1 个 batch 且包含全部 129 条——证明原子驱动不分块splits writes into KV-sized batches在普通RecordingDriver下同样的 129 条被切成[128, 1]两个 batchsplits writes by KV batch payload size9 条各 120 KiB 的 name总计约 1080 KiB 超过 976 KiB 上限被拆成多个 batch每个都满足条数与字节数约束且 reload 后 nameRegistry 完全一致。三、dirty 标记的清除时机先成功后清除3.1 为什么必须在所有操作成功后清除CLAUDE.md 强调的clears dirty markers only after all write/delete operations succeed是可靠性关键。flush()中所有dirty false的赋值都发生在 batch 写入driver.batch/ 分块 batch和runDeletes全部成功返回之后storage.ts。其正确性价值在于失败的 flush 不会丢数据如果某次 batch 写入中途失败dirty 标记保持true内存中的数据与磁盘KV状态不一致但下次flush()会原样重试把上次没写进去的数据重新写入避免已清除但实际没落盘的窗口如果先清除 dirty 再写一旦写失败内存就认为已持久化崩溃后重放将缺失历史导致工作流状态回退甚至 HistoryDivergedError。3.2 测试如何验证失败不清除storage.test.ts 的keeps dirty markers when a batch write fails用例专门验证了这一契约构造一个 step 条目 completed 元数据用failOnBatch 1让驱动在第一次 batch 时抛injected batch failure断言flush()抛出该错误且entry.dirty true、metadata.dirty true、flushedNameCount 0清除failOnBatch后再次flush()才看到dirty变为false。这从测试层面锁定了清除时机的行为——它不是实现细节而是被测试强制保证的对外契约。四、删除操作的批处理PendingDeletions 与并发轮次flush 不仅要写还要删。工作流循环loop在迭代中会按historyPruneInterval裁剪旧历史裁剪收集到的删除会以PendingDeletions的形式与下一次 state 写入一起提交export interface PendingDeletions { prefixes: Uint8Array[]; keys: Uint8Array[]; ranges: { start: Uint8Array; end: Uint8Array }[]; }见 storage.tsrunDeletesstorage.ts把三种删除操作归一为一批操作driver.deletePrefix(prefix)按前缀整体删除如循环某个 location 前缀的历史driver.deleteRange(start, end)按半开区间删除driver.batchDelete(chunk)按键删除先用splitBatchDeletes按MAX_KV_BATCH_ENTRIES128切块再每个 chunk 一次事务。然后按MAX_CONCURRENT_DELETES64为步长分轮次执行Promise.all即每轮最多 64 个删除操作并发保证不会突破 KV 事务配额。runDeletes返回是否有任何删除发生有则标记historyUpdated并触发onHistoryUpdated()通知。collectDeletionsForPrefixstorage.ts则负责从内存中同步移除匹配前缀的条目与元数据并生成驱动级删除操作——它可被立即执行也可与下一次 flush 一起延迟执行deleteEntriesWithPrefix提供立即执行版本。五、持久化 Schema从 BARE 定义到序列化栈CLAUDE.md 的第二条指引是关于持久化 schema 的双端同步。workflow engine 的持久化格式不是手写的结构体而是由一份 BAREBinary Application Record Encodingschema 定义并代码生成。5.1 schema 定义了哪些类型schemas/v1.bare 完整定义了持久化层的数据契约Location 体系NameIndexu32指向名字注册表、LoopIterationMarkerloop iteration、PathSegmentunion、LocationlistPathSegment。这正是 architecture.md 中路径式 location NameIndex 优化的二进制形态——用数值索引代替字符串减小重复命名的存储开销Entry 体系StepEntryoutput CBOR error、LoopEntrystate CBOR iteration output、SleepEntrydeadline SleepState、MessageEntry、RollbackCheckpointEntry、JoinEntrybranches map、RaceEntrywinner branches、RemovedEntry迁移占位、VersionCheckEntry由EntryKindunion 归并为Entryid location kind元数据EntryMetadatastatus、attempts、lastAttemptAt、createdAt、completedAt、rollbackCompletedAt、rollbackError工作流级WorkflowState枚举PENDING/RUNNING/SLEEPING/FAILED/COMPLETED/ROLLING_BACK、WorkflowMetadatastate output error version 哈希。其中所有用户数据Cbor类型都用 CBOR 编码配合cbor-x运行时依赖见 package.json。5.2 编译与序列化栈Schema 到运行时的流水线是BARE → TypeScript由compile:bare脚本执行tsx scripts/compile-bare.ts compile schemas/v1.bare -o dist/schemas/v1.tspackage.json底层使用bare-ts/tools的transform()compile-bare.ts。后处理生成代码做两处替换——把bare-ts/lib导入换成rivetkit/bare-ts项目内维护的 fork并移除 Node.js 的assert导入、注入一个纯 JS 的自定义assert函数以保证生成代码可在非 Node 环境如浏览器 inspector运行。脚本注释明确要求与engine/packages/runner-protocol/build.rs保持同步compile-bare.ts。版本化schemas/versioned.ts用vbare的createVersionedDataHandler包装 v1 生成的编解码器CURRENT_VERSION 1为将来格式演进预留版本通道。内部类型桥接schemas/serde.ts把引擎内部 TypeScript 类型与 BARE 生成类型互转并将任意用户值用 CBOR 编码为Cbor字段。5.3 为什么需要双端同步CLAUDE.md 明确指出The workflow engine persistence schema is duplicated in RivetKit for inspector transport.即这份持久化 schema在 RivetKit 主包中被复制了一份用于 Inspector调试器/观察者传输——Inspector 需要在不执行工作流的情况下解码其历史、输出与状态参见 rivetkit 的 inspector 模块 对 workflow 状态的实验性暴露。这意味着同一份二进制格式存在两个实现入口workflow-engine 侧执行引擎读写自身持久化rivetkit 侧Inspector 传输解码。任何一侧改动格式新增 Entry 类型、修改字段、调整 union 标签而不同步另一侧都会导致 Inspector 解码失败或工作流历史被误读。这就是镜像必须存在的原因。六、Schema 更新与重建一份可执行的同步清单CLAUDE.md 给出了修改持久化 schema 时的完整操作步骤将其整理为可执行的清单第 1 步同时修改两份 schema 源文件主定义rivetkit-typescript/packages/workflow-engine/schemas/v1.bare镜像定义文档指示的镜像路径为rivetkit-typescript/packages/rivetkit/schemas/persist/v1.bare第 2 步重新编译两端生成代码# 编译 workflow-engine 的 schema 为 TypeScript生成 dist/schemas/v1.ts pnpm -C rivetkit-typescript/packages/workflow-engine run compile:bare # 重新构建 rivetkit含 Inspector 传输侧 schema pnpm -C rivetkit-typescript/packages/rivetkit run build:schema第 3 步验证运行 workflow-engine 的测试套件确认编解码与回放不回归pnpm -C rivetkit-typescript/packages/workflow-engine test关注 schema 相关测试如 version.test.ts以确认版本化处理正常。需要说明的是截至本仓库当前状态rivetkit-typescript/packages/rivetkit目录下实际存在的生成文件为src/common/bare/actor-persist/下的v1.ts至v4.ts见 actor-persist-versioned.ts而 CLAUDE.md 所写的schemas/persist/v1.bare镜像源文件在当前工作区中尚未出现。因此实际操作时应以各自包内实际的 schema 脚本与生成位置为准先确认 rivetkit 侧镜像文件的确切路径再修改避免在错误的目录写入。这条指引的核心纪律不变——主定义与镜像必须同改、同步构建、同步验证。七、实践要点与排查建议7.1 何时会触发 976 KiB 上限最典型的触发场景是单个 entry 或单条写入过大。例如 step 输出一个大型对象序列化为 CBOR 后超过 976 KiBflush 时splitBatchWrites会直接抛错。排查路径定位是单条超限还是批量超限——单条超限错误信息会直接包含字节数与上限检查对应 step 的输出是否包含过大载荷考虑拆分为多个 step 或改用 KV 直接存储大对象若是批量超限多步并发产出大量小写入让驱动声明atomicBatch交给后端原子处理或调整 step 粒度。7.2 删除风暴与并发上限一次循环历史裁剪prune可能产生大量删除操作。runDeletes以 64 为并发轮次边界防止瞬间打满 KV 事务配额同时裁剪与状态写入在同一逻辑 flush 内完成保证新状态 旧历史清理的一致性。7.3 不要手动清除 dirty任何对 Storage 的扩展或自定义 Driver 集成都不应在flush()之外手动将dirty置为false——dirty 标记是内存与持久层一致的唯一凭证其清除时机被上述测试强制锁定storage.test.ts违反它将直接破坏崩溃恢复的正确性。7.4 改 schema 时的纪律永远同时改 workflow-engine 主 schema 与 rivetkit 镜像重建后跑两端构建与测试重点验证 Inspector 传输解码为已部署的历史数据考虑版本化vbare 的createVersionedDataHandler是向后兼容的入口不要直接原地修改既有字段类型。八、结语storage.flush(...)是 RivetKit Workflow Engine 持久化正确性的最后一公里128 条 / 976 KiB 的批处理约束让每次写入都符合 Actor KV 的事务上限atomicBatch让引擎在不同后端的原子性承诺下都能正确工作而先成功、后清除 dirty的时机则保证了任何一次崩溃都不会造成历史丢失。在此基础上BARE schema 双端镜像同步确保了执行引擎与 Inspector 对同一份二进制格式的理解始终一致。理解这两条指引就掌握了该引擎存储层的核心契约也为排查写入超限、扩展持久化格式打下了基础。进一步的实现细节可继续阅读 architecture.md存储键编码、Location 系统、隔离模型与 QUICKSTART.md引擎 API 使用示例。【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考