- 数据库
- 流处理
- 后端
- 数据工程
【免费下载链接】risingwave
Event streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.
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 对外提供五个核心 API:get_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.rs(StorageTable及get_row、按扫描范围构造的new/new_partial)。
写路径(Write Path):epoch 与 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_insert(test_state_table.rs)先构造StateTable并init_epoch,随后连续insert、delete,再通过get_row验证内存态与存储态合并后的读取结果——这恰好覆盖了「写入经 MemTable 缓存、读取需合并」的核心行为。
读路径(Read Path):MemTable 与 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]要点在于:pk=2 的行先插入后被删除,因此读不到;pk=3 的行在 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 State:Sum、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 State:Max、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值;- 该设计允许流式算子在缓存 miss(cache 失效)时不必全量读取存储数据,而只需读取一部分(排序后最靠前的部分即可恢复极值);
- 由于流中可能含有
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 提供关系语义 API,Mem Table 缓存未提交操作,Storage Table 提供只读的部分列输出;读路径通过
StateTableIter合并内存与存储,保证"未提交数据可见",写路径由 barrier 驱动、按 epoch 批量 flush 到 Hummock; - Schema 随需求演化:聚合内部表从值状态的最小 Schema(
table_id/group_key),演化到极值状态追加sort_key/upstream_pk以支持增量恢复,再到 append-only 标志下回归最小 Schema——每一列都有明确的正确性动机。
如果想深入底层存储,建议继续阅读 state-store-overview.md 了解 Hummock 的 epoch/version/checkpoint 语义;想了解行与 KV 字节层面的编码细节,可阅读>赞
- 数据库
- 流处理
- 后端
- 数据工程
【免费下载链接】risingwave
Event streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.
相关推荐
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组件设计系统