设计解析:流式状态如何以行式编码写入 KV 状态存储)
数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载RisingWave 采用关系型数据模型所有流式计算算子executor的有状态数据都持久化到以 Hummock 为底层引擎的 KV 状态存储中。本文以 relational-table.md 为核心深入讲解 RisingWave 如何通过「关系表层」Relational Table Layer在 KV 状态存储之上提供行式关系语义的读写接口包括行式编码Row-based Encoding的选型理由、State Table / MemTable / Storage Table 三层组件的职责划分、基于 epoch-barrier 机制的写路径与合并读路径并以 HashAgg 的sum/count与max/min两种内部表结构为例给出完整的表 Schema 设计。读完本文你将理解流式算子状态在 RisingWave 中的实际存储形态以及 append-only 标志如何显著简化内部表的 Schema 设计。背景从关系模型到 KV 存储RisingWave 采用关系型数据模型。关系表包括用户表 table 和物化视图 materialized view由一组命名且强类型的列构成数据模型与编码细节可参考 />组件使用模式职责State Table流式模式读写为流式算子提供关系语义的读写接口Mem Table流式模式内存缓冲缓存一个 epoch 内的表操作是未提交数据的暂存区Storage Table批式模式只读按上层需要的部分列输出存储中的数据State Table 对外提供五个核心 APIget_row、scan、insert_row、delete_row、update_row分别是流式算子的点查、范围扫描与增删改接口。这些接口在源码中均有对应实现例如 state_table.rs 的pub async fn get_row、state_table.rs 的insert/delete/update以及 state_table.rs 中按 vnode 与 pk 范围扫描的iter_with_vnode。Mem Table 是内存中的缓冲结构缓存一个 epoch 内的所有操作Storage Table 则只读并且可以只输出上层需要的部分列partial columns。对应批式端的实现可参考 batch_table/mod.rsStorageTable及get_row、按扫描范围构造的new/new_partial。写路径Write Pathepoch 与 barrier 驱动的两阶段写入文档给出了清晰的写流程流式算子调用 State Table 的 API 执行操作如insert、delete这些操作先缓存在 Mem Table 中当一条 barrier屏障流过该算子时算子把缓存的操作flush 到 KV 状态存储flush 时State Table 把缓存操作按行式格式转换为 KV 对并携带特定的 epoch写入 Hummock。以文档示例说明算子依次执行insert(a, b, c)与delete(d, e, f)Mem Table 先把这两个操作缓存在内存收到新的 barrier 后State Table 将二者按行式编码转换为 KV 操作写入 Hummock。写入过程中 KV 对都会关联到 barrier 对应的 epoch——这正是 barrier 检查点算法在存储层的落地每个 epoch 的流式状态被整体提交。底层写入以「排序且 key 唯一的批量」形式调用 Hummock 的ingest_batch这一点在 state-store-overview.md 中有详细约束说明。文档配套的写路径示意图见 relational-table-03.svg在单元测试中也能看到同样的流程test_state_table_update_inserttest_state_table.rs先构造StateTable并init_epoch随后连续insert、delete再通过get_row验证内存态与存储态合并后的读取结果——这恰好覆盖了「写入经 MemTable 缓存、读取需合并」的核心行为。读路径Read PathMemTable 与 Storage 的合并流式算子需要读到最新写入、尚未提交uncommitted的数据因此共享存储state store中的数据不是最新的——Mem Table内存中的数据总是比共享存储中的数据更新。State Table 的get与scan因此都必须合并 Mem Table 与 Storage Table 的数据后再返回结果。Get点查示例假设关系表第一列是主键 pk依次执行以下操作insert [1, 11, 111] insert [2, 22, 222] delete [2, 22, 222] insert [3, 33, 333] commit insert [3, 3333, 3333]commit 之后又插入一条新记录则各 pk 的Get结果为Get(pk 1): [1, 11, 111] Get(pk 2): None Get(pk 3): [3, 3333, 3333]要点在于pk2 的行先插入后被删除因此读不到pk3 的行在 commit 前插入的是[3, 33, 333]commit 后又更新为[3, 3333, 3333]由于未提交的新数据在 Mem Table 中比存储更新点查返回的是内存中的[3, 3333, 3333]。Scan范围扫描StateTableIter 合并迭代器scan由StateTableIter实现它是MemTableIter读 Mem Table与StorageIter读 KV 存储的合并迭代器merge iterator。合并规则是若某个 pk 同时存在于共享存储与 MemTable则返回MemTableIter的结果——因为内存中的数据更新。文档中的示例relational-table-02.svg展示了StateTableIter按顺序产出1 - 4 - 5 - 6的扫描结果。该设计与 Hummock 自身的读取模型一致Hummock 读取时也是用MergeIterator合并多个 SST 后按 epoch 选出可见版本详见 state-store-overview.md。关系表层在更高一层、以行为粒度实现了类似的新数据优先合并语义并通过iter_with_vnode等接口限定 vnode 与主键范围以控制扫描粒度。实战案例HashAgg 的内部表 Schema 设计下面以 HashAgg 的两类典型状态——值状态value state如sum、count与极值状态extreme state如max、min为例展示内部表的 Schema 是如何随需求演化的。相关的优化器实现可参考 logical_agg.rs 中对聚合调用的规划逻辑。table_id全局唯一的关系表标识table_id是元数据服务meta为每个关系表对象分配的全局唯一 ID。Meta 在收到查询计划后负责遍历 Plan Tree计算所需关系表的总数并分配 ID。例如Hash Join 算子需要2张表左表 1 张、右表 1 张Agg 算子需要的表数量取决于聚合调用的个数agg call 数量。值状态Value StateSum、Count查询示例select sum(v2), count(v3) from t group by v1该查询需要初始化2 张关系表sum(v2)一张、count(v3)一张表 Schema 为table_id / group_key即以table_id隔离表以group_key分组键v1作为键。由于值状态是可增量的新数据只需在旧聚合值上累加每张表只需按分组键存一个聚合值无需保留全部明细数据。极值状态Extreme StateMax、Min查询示例select max(v2), min(v3) from t group by v1该查询同样需要初始化2 张关系表。当上游不是 append-only 时表 Schema 变为table_id / group_key / sort_key / upstream_pk与值状态的关键差异在于多出了sort_key与upstream_pk设计动机如下sort_key的排序方向取决于聚合函数类型max()的sort_key按Ascending升序排列min()的sort_key按Descending降序排列。这样可以把极值行放在扫描顺序的最前面方便快速取到当前最大/最小值upstream_pk上游主键被附加到键尾用于保证键的唯一性——因为不同上游行可能携带相同的sort_key值该设计允许流式算子在缓存 misscache 失效时不必全量读取存储数据而只需读取一部分排序后最靠前的部分即可恢复极值由于流中可能含有update或delete操作不可能在只存单个值的情况下始终得到正确结果因此算子会把所有流式数据尽量写入存储——这正是极值状态表需要保留全部明细而非只存一个值的根本原因。append-only 优化Schema 的显著简化如果t以 append-only 标志创建极值状态的表 Schema 变为table_id / group_key与值状态完全一致。原因在于append-only 模式下不存在update或delete操作因此缓存永远不会 miss缓存里缺失的数据不可能被删改我们只需向存储写入一个值即可无需再保存sort_key与upstream_pk来支撑极值恢复。这一设计体现了 RisingWave 把流形态append-only 与否作为 Schema 化简关键输入的思路。小结三个关键设计决策回到本文开头的问题可以把本设计归纳为三个相互咬合的决策行式编码由于流式状态总是整行读写、无部分更新需求以「一行 一个 KV 对」的方式存储用更少的 KV 对换更高的读写性能三层读写分层State Table 提供关系语义 APIMem Table 缓存未提交操作Storage Table 提供只读的部分列输出读路径通过StateTableIter合并内存与存储保证未提交数据可见写路径由 barrier 驱动、按 epoch 批量 flush 到 HummockSchema 随需求演化聚合内部表从值状态的最小 Schematable_id/group_key演化到极值状态追加sort_key/upstream_pk以支持增量恢复再到 append-only 标志下回归最小 Schema——每一列都有明确的正确性动机。如果想深入底层存储建议继续阅读 state-store-overview.md 了解 Hummock 的 epoch/version/checkpoint 语义想了解行与 KV 字节层面的编码细节可阅读>赞分享数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载相关推荐Rivet Actors 引擎的嵌入式 KV 选型基于 RocksDB 的 Actor 状态存储设计决策Rivet Actors 引擎的嵌入式 KV 选型基于 RocksDB 的 Actor 状态存储设计决策 本文解析 Rivet Actors 引擎中“嵌入式后端AI Agent人工智能流程编排WebSocketNomad 状态存储架构深入解析基于 Raft 的 FSM 与内存状态库设计Nomad 状态存储架构深入解析基于 Raft 的 FSM 与内存状态库设计 导读 本文围绕 Nomad 服务器端最核心的架构组件——状态存储State S任务调度云原生运维后端Radix Vue状态提升模式何时以及如何共享组件状态Radix Vue状态提升模式何时以及如何共享组件状态 你是否遇到过Vue组件间状态同步的难题当多个复选框需要共享选中状态、表单元素需要跨组件联动时传统的前端UI组件设计系统上一篇Thorium-Win性能调优10个简单设置让浏览器飞起来下一篇Ungit社区与支持如何参与贡献并解决使用难题创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考