数据库流处理后端数据工程【免费下载链接】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 的流式计算与存储层围绕两类 Key 展开设计用于在流中唯一标识记录、驱动下游状态维护的Stream Key以及用于决定存储层数据排序与检索效率的存储主键Storage Primary Key。本文以仓库设计文档 docs/dev/src/design/keys.md 为骨架结合前端优化器与编目层源码系统讲解两类 Key 的定义、约束、推导逻辑与协同关系帮助读者理解EXPLAIN计划中stream_key、pk_columns、pk_conflict等字段的真实含义并掌握如何在物化视图、表与 Sink 的建表建视图实践中正确使用 ORDER BY、DISTRIBUTE BY 等语法来影响 Key 的构成。一、两类 Key 的总体定位在 RisingWave 中Key 并不只有 SQL 语义下的主键PRIMARY KEY一种。流式执行与存储分别对记录标识提出了不同的要求最终演化出两类职责互补的 KeyKey 名称作用域核心职责Stream Key流Stream在流中唯一标识一条记录供下游算子维护 per-record 状态、实现更新Update/删除Delete对齐存储主键Storage Primary Key内部又称 pk存储层唯一标识存储中的记录并决定 Key-Value 在存储中的排序从而支持有序扫描两者的关系可以概括为Stream Key 负责流上如何定位记录存储主键负责存储中如何组织记录而更新流恰好通过 Stream Key 找到存储主键对应的整条记录再完成状态的定点更新。二、Stream Key流内记录的身份标识2.1 定义与直观示例Stream Key 是能在 RisingWave 流中唯一标识一条记录的列组合。设计文档给出了一个非常直观的示例假设某个流分块Stream Chunk的 Stream Key 为k1, k2| op | k1 | k2 | v1 | v2 | |----|----|----|----|----| | - | 1 | 2 | 1 | 1 | | | 1 | 2 | 3 | 4 | | | 0 | 1 | 2 | 3 |其中op列表示操作类型-表示删除Delete表示插入Insert。对照 Stream Key(k1, k2)可以读出对于键(1, 2)记录从(1, 2, 1, 1)更新为(1, 2, 3, 4)即先删除旧值、再插入新值对于键(0, 1)记录(0, 1, 2, 3)是新增插入。2.2 Stream Key 不要求是最小标识集合文档特别强调Stream Key不一定是能够标识记录的最小列集合。最典型的情况是分组键group key即分布键/distribution key也会被并入 Stream Key用于规定记录的分布方式。这在源码中得到印证前端优化器的 derive.rs 中的 derive_pk 函数 在推导 Stream Key 时首先将用户要求的分布方式user_distributed_by中的分布列dist_column_indices加入 Stream Key然后再追加输入计划的expect_stream_key()并做去重。也就是说只要一条 SQL 指定了DISTRIBUTE BY之类影响分布的要求对应的列就会成为 Stream Key 的一部分——即便仅凭这些列并不能唯一标识记录但它们保证了相同键的记录一定落在同一个计算节点/分区上这是流式 join、聚合等算子本地维护状态的前提。2.3 流上的一致性约束Insert 与 Delete 必须交替既然下游算子依赖 Stream Key 来维护每个键的独立状态流本身就必须满足以下约束文档原文语义同一个 Stream Key不允许连续出现两次 Insert中间必须夹着一次 Delete同一个 Stream Key不允许连续出现两次 Delete中间必须夹着一次 Insert对于更新UpdateDelete 侧的值必须匹配该 Stream Key 上一次产出的旧值。举例若 Stream Key(1, 2)的上一次值是(1, 2, 1, 1)那么更新到(1, 2, 3, 4)必须表示为| op | k1 | k2 | v1 | v2 | |----|----|----|----|----| | - | 1 | 2 | 1 | 1 | | | 1 | 2 | 3 | 4 |删除一个与旧值不同的记录是非法操作。原因正如文档所述下游算子正是以 Stream Key 为索引来维护各自的 per-record 状态聚合中间值、Join 缓存、Top-N 堆等如果 Delete 的旧值与算子状态中的记录不一致状态将无法正确回退最终导致结果错误或状态泄漏。2.4 源码中的体现Stream Key 在代码中直接落为计划节点的属性。例如physical_table.rs 中TableFragments相关的表目录以stream_key: Vecusize记录 Stream Key 对应的列索引usize指向输出列的序号external_table.rs 中外部表同样维护stream_key: Vecusize并通过 protobuf 序列化传递StreamMaterialize节点在 stream_materialize.rs 构造时以Some(table.stream_key())写入自身的 PlanBase使下游能直接读取本节点产出的 Stream Key。此外TTL数据过期场景下系统会把 TTL watermark 列追加到 Stream Key 中。相关逻辑位于 stream_materialize.rs 的 derive_table_catalog当一行数据进入带 TTL 的表却查不到时无法判断它是新行还是对已过期行的更新把 TTL watermark 列加入 Stream Key 可以确保下游作业视角下不会出现双重插入。三、存储主键Storage Primary Key存储层的有序骨架3.1 与 SQL 主键的本质区别文档明确指出这里讨论的Primary Key是流式算子中常见的内部主键pk它不同于 SQL 里的 PRIMARY KEY更恰当的名字是Storage Primary Key存储主键。它的职责有两层唯一标识存储中的一条记录——类似传统数据库主键提供排序属性ordering property——这是它在 RisingWave 存储模型中更关键的作用。RisingWave 的存储层Hummock本质是一个按 Key 排序的 Key-Value 存储扫描结果天然按 Key 有序。因此物化视图的状态表如何排序完全由存储主键决定。3.2 一个完整的例子ORDER BY 如何改写存储主键文档给出了经典示例create table t1(id bigint primary key, i bigint); create materialized view mv1 as select id, i from t1 order by i, id;mv1的执行计划EXPLAIN输出如下StreamMaterialize { columns: [id, i], stream_key: [id], pk_columns: [i, id], -- notice the pk_columns pk_conflict: NoCheck } └─StreamTableScan { table: t1, columns: [id, i] }注意虽然上游表的 SQL 主键只是id但pk_columns却是[i, id]。原因正是存储层 Key-Value按键排序的性质物化视图声明了order by i, id为了在遍历存储时能按正确的顺序逐条产出记录存储主键必须为[i, id]其中i位于键的前缀位置。这样按分区per partition顺序迭代 Key 时返回的记录天然满足i, id的排序无需额外排序开销。3.3 更新流如何与存储主键协作存储主键除了排序还承担定位功能。当更新流到达时系统先利用Stream Key即id定位需要更新的记录通过该id取出存储中的整条旧记录从旧记录中读出i的旧值用(i, id)组成存储主键对物化状态进行定点更新。这正是两类 Key 协同的典型链路Stream Key 负责流上的身份存储主键负责存储中的地址二者通过取整条记录、再读前缀列的方式完成桥接。文档强调不能删除同一个 Stream Key 下不同的旧值也正是因为下游算子含物化状态依赖这条链路做状态维护。3.4 源码级的推导逻辑derive_pk存储主键与 Stream Key 的推导集中在 derive.rs 的 derive_pk该函数同时服务于表和 Sink 的建表/建 Sink 流程。其核心逻辑可归纳为四步确定 Stream Key 初值取分布键dist_column_indices若无分布要求则为空并入输入计划的 Stream Key将input.expect_stream_key()去重后追加保证流上可标识用用户 ORDER BY 构造存储主键前缀对user_order_by做函数依赖最小化func_dep.minimize_order_key后逐列加入pk同时维护remaining_stream_key一旦 ORDER BY 列已覆盖全部 Stream Key 列就提前停止对应源码中的stop_order_by_after_stream_key补齐剩余 Stream Key 列将尚未进入pk的 Stream Key 列以升序OrderType::ascending追加到pk尾部。由此可以推出一个重要不变量Stream Key 一定是存储主键的子集stream_key ⊆ pk_columns且存储主键的前缀由用户 ORDER BY 决定、尾部由 Stream Key 补齐——这既保证了排序要求又保证了按 Stream Key 定位后一定能构造出完整的存储键。另外create table场景下pk_column_indices直接来自用户声明的 SQL 主键Stream Key 与表主键一致stream_materialize.rs。3.5pk_conflict写入冲突策略EXPLAIN中的pk_conflict字段对应表目录中的ConflictBehavior它决定了当写入记录与存储主键冲突时的处理方式。从 stream_materialize.rs 可以看到NoCheck默认策略适用于 Append-Only / Retract 流此时系统会拒绝 Upsert 输入reject_upsert_input!保证不会出现主键冲突写入Overwrite、IgnoreConflict、DoUpdateIfNotNull用于 Upsert 流可将 Upsert 流转换为 Retract 流语义后继续处理。在EXPLAIN输出中columns / stream_key / pk_columns / pk_conflict / watermark_columns正是由 stream_materialize.rs 的 distill 实现 逐字段渲染出来的其中pk_columns直接取自table.pk中每个ColumnOrder的列名stream_key取自table.stream_key()对应的列名。四、两类 Key 的协同全景把前面内容串起来RisingWave 中一条记录的生命周期大致是写入/更新到达外部变更进入流携带 Stream Key流式算子维护状态聚合、Join、Top-N 等算子按 Stream Key 索引各自状态依靠 Insert/Delete 交替约束正确回退物化落盘StreamMaterialize节点按存储主键ORDER BY 前缀 Stream Key 补齐把记录写入 Hummock有序读取按分区顺序扫描存储键直接得到满足 ORDER BY 的结果定点更新用 Stream Key 命中记录、取出完整旧值、以存储主键定位并覆盖。这一设计同时满足了三个诉求流上状态一致Stream Key 交替约束、存储上有序存储主键前缀排序、分布均匀可扩展分布键并入 Stream Key。对于下游算子而言Stream Key 是状态寻址的句柄对于存储引擎而言存储主键是排序与定位的目录。五、实践要点速查在EXPLAIN/EXPLAIN (FORMAT JSON|YAML)中查看StreamMaterialize节点时优先核对stream_key与pk_columns若两者不一致说明存在 ORDER BY 参与存储主键、或分布键并入 Stream Key 的情况物化视图的ORDER BY会改变存储主键从而影响状态表在 Hummock 中的物理排序与扫描效率应把最常用的排序/过滤前缀列放在 ORDER BY 前部DISTRIBUTE BY分布键会进入 Stream Key但通常不会成为唯一标识列它影响的是记录的分布亲和性流式计划要求同一 Stream Key 的 Insert/Delete 严格交替且更新必须携带正确旧值任何破坏该约束的上游如非法定制连接器都会导致下游状态错乱带 TTL 的表会自动把 TTL watermark 列追加进 Stream Key以避免过期数据被误判为新插入。对 Key 语义的深入理解能帮助你在排查流式结果错乱、状态膨胀、扫描性能异常等典型问题时第一时间从stream_key/pk_columns中找到根源。更多细节可继续阅读仓库中的设计文档 docs/dev/src/design/keys.md 与前端优化器实现 derive.rs、stream_materialize.rs。赞分享数据库流处理后端数据工程【免费下载链接】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点击查看免费下载相关推荐Matter SDK 示例应用深度解析Qorvo QPG6200 Persistent Storage 应用与 Key-Value 存储 API 验证Matter SDK 示例应用深度解析Qorvo QPG6200 Persistent Storage 应用与 Key Value 存储 API 验证 本篇文物联网智能家居嵌入式通信Milvus 按主键检索Search By Primary Keys / Search by IDs设计与实现深度解析Milvus 按主键检索Search By Primary Keys / Search by IDs设计与实现深度解析 在 Milvus 向量数据库中标准数据库向量数据库分布式数据库后端ToolJet Database 主键Primary Key完整指南单字段、复合主键的创建、修改与删除ToolJet Database 主键Primary Key完整指南单字段、复合主键的创建、修改与删除 ToolJet Database 是 ToolJe低代码后端前端AI 应用MCP 服务上一篇 Integration Task下一篇Aqua构建现代网站和用户系统的开源利器创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
