RFC 深度解读:基于序列(Sequence)与分区域水位线的正确性优先增量计算方案)
GreptimeDB Flow 增量查询Incremental QueryRFC 深度解读基于序列Sequence与分区域水位线的正确性优先增量计算方案【免费下载链接】greptimedbThe open-source observability database. One columnar engine for metrics, logs, and traces, on object storage.项目地址: https://gitcode.com/GitHub_Trending/gr/greptimedb导读本文以 docs/rfcs/2026-03-16-flow-inc-query.md 为骨架深度解读 GreptimeDB Flow 引擎引入的基于序列Sequence的增量查询计划Lite 版。这份 RFC 提出的核心思路是让 Flow 批处理查询只读取seq checkpoint的新增数据并借助分区域正确性水位线per-region correctness watermark推进检查点从而避免在同一时间窗口内反复重算旧数据当增量读取过期stale或正确性无法证明时则确定性回退到全量重算。读完本文你将理解三个 Flow 查询扩展键的含义与用法、IncrementalQueryStale过期保护机制、memtable→SST 冲刷边界的正确性局限以及 FullSnapshot ↔ Incremental 双状态机的完整迁移与回退策略并能在仓库源码层面找到每一处设计的落地证据。一、背景与动机为什么 Flow 批处理需要增量查询GreptimeDB 的 Flow 引擎以批处理batching方式周期性地对源表执行查询并将聚合结果写入 Sink 表。在早期实现中即使某个时间窗口的数据大部分已经计算过Flow 依然会反复重算同一窗口内的旧数据。这种全量重算在数据持续写入、窗口不断滚动的高频场景下会造成明显的计算浪费。RFC 的出发点非常明确复用已有的序列扫描sequence scan能力——存储引擎本身就能按行级序列号过滤数据——让 Flow 在下一轮批处理中只读取自上次检查点以来新增的行seq checkpoint从而显著降低重复计算的开销。但这里存在一个关键的工程张力性能优化必须以正确性为前提。因此本 RFC 通篇的基调是correctness first正确性优先只有当增量读取的正确性可以被严格证明时系统才允许走增量路径任何无法证明的场景都无条件回退到全量重算。本 RFC 的目标与非目标RFC 明确了四个目标Goals为 Flow 引入可选开启的增量读取能力seq given_seq返回分区域正确性水位线用于推进检查点除非显式开启保持既有查询行为完全不变默认兼容为过期stale或无法证明正确性的增量读取定义确定性回退策略。同时划定了三个非目标Non-Goals避免范围蔓延不改变业务 schema——不在结果行中引入合成的 watermark 列防止 schema 污染v1 不做全局吞吐优化——先保证正确性性能优化留到 Phase 2正确性无法证明时不输出观测性水位线——宁可少给信息也不给错误信息。二、核心设计一三个查询扩展键Query Options增量能力通过QueryContext扩展extension键对外暴露均为可选开启opt-in且只影响 Flow 的增量执行路径。RFC 定义三个键源码中它们的常量定义位于 src/query/src/options.rs扩展键类型含义flow.incremental_after_seqsJSON 字符串region_id - seq映射每个源 Region 的排他性序列下界增量读取只返回seq given_seq的行flow.incremental_mode字符串枚举增量读取模式当前唯一合法值是memtable_onlyflow.return_region_seq布尔字符串true/false调用方是否期望在终端 metrics 中收到分区域水位线元数据代码中的完整定义src/query/src/options.rs还包含第四个键flow.sink_table_id它用于区分源表扫描与 Sink 表读取——Sink 表的扫描应像普通读一样不参与增量下界与快照绑定逻辑。2.1 解析与校验parse_flow_extensionsFlowQueryExtensions::parse_flow_extensions 负责把扩展字符串解析为结构体其规则与原文档完全对应且更细致若没有任何 Flow 相关键返回Ok(None)普通查询路径完全不受影响默认兼容目标flow.incremental_mode的值需与memtable_only大小写不敏感匹配否则报Invalid value for flow.incremental_mode错误flow.incremental_after_seqs需为合法 JSONregion_id必须是可解析为u64的字符串键序列值可以是数字或字符串形式例如{1:10,2:20}与{1:10,2:20}均合法源码中有对应单测 test_parse_flow_extensions_after_seqs_string_values值为布尔等其他类型会被拒绝flow.return_region_seq必须是true/false大小写不敏感约束规则当flow.incremental_modememtable_only时flow.incremental_after_seqs为必填且不能为空——RFC 原文若必需参数缺失或非法返回参数错误在代码中被精确落实src/query/src/options.rs。2.2 扫描级校验validate_for_scan对每个源 Regionvalidate_for_scan 还会做二次校验若请求是memtable_only增量模式则incremental_after_seqs中必须包含该源 Region 的条目否则返回Missing region id in flow.incremental_after_seqs错误同时若扫描目标是sink_table_id对应的表则跳过增量语义返回false。这一先整体解析、再逐 Region 校验的两段式设计保证了多 Region 查询中任何一个源 Region 缺少下界时查询会被拒绝而不是静默退化为全量读取。2.3 客户端如何携带这些扩展仓库客户端层已经提供了携带扩展的完整示例见 src/client/src/database.rs(flow.return_region_seq, true), (flow.incremental_after_seqs, r#{1:10,2:20}#),而 Flow 任务侧则在 src/flow/src/batching_mode/task/inc.rs 的build_flow_query_extensions中统一构造增量扩展始终附加flow.return_region_seqtrue当处于增量模式且持有非空检查点时再附加flow.sink_table_id、flow.incremental_modememtable_only以及序列化后的检查点映射flow.incremental_after_seqs。三、核心设计二扫描映射与 memtable→SST 边界的正确性局限3.1 下界与上界的映射当增量模式开启时扫描请求的序列边界按如下规则映射RFC 第 2 节实现位于 src/query/src/dummy_catalog.rs 的decide_flow_scan下界flow.incremental_after_seqs中该 Region 的值映射为memtable_min_sequence即排他性exclusive下界——只读取序列严格大于该值的行上界保留既有的快照行为即memtable_max_sequence来自查询上下文已绑定的分区域快照序列query_ctx.get_snapshot(region_id)skip_sst_files仅当同时满足增量模式为memtable_only且下界存在时才置位——即 memtable-only 增量 delta 只扫 memtable、跳过 SST 文件而仅有快照上界、无增量下界的带围栏的全量快照读则必须保留 SST以便存储引擎能对上界过期做拒绝。最终ScanRequest通过 build_scan_request 组装其中还包含sst_min_sequence、snapshot_on_scan等字段。若必需的增量参数缺失或非法validate_for_scan会直接返回参数错误Argument Error不会带着残缺参数继续执行。3.2 v1 的关键限制SST 缺少行级序列元数据RFC 明确指出 v1 最重要的正确性局限增量过滤仅在 memtable 行上被证明正确——memtable 保留行级序列号可以精确执行seq checkpointSST 文件不保留详细的行级序列元数据只暴露更粗粒度的文件级序列信息因此seq checkpoint不能假设能跨 memtable→SST 冲刷边界做精确的增量裁剪。换言之一旦某行数据从 memtable 刷入 SST系统就无法仅凭seq checkpoint证明该行已包含在上次计算结果中此时继续做尽力而为的增量扫描是危险的。正确的行为是回退到全量重算见下文第五节而不是带着不精确的增量继续。四、核心设计三Stale 保护与IncrementalQueryStale错误4.1 错误定义与三分支判定RFC 新增专用过期错误IncrementalQueryStale { region_id, given_seq, min_readable_seq }判定语义非常清晰条件结果given_seq min_readable_seq返回IncrementalQueryStale错误given_seq min_readable_seq查询有效读取seq given_seqgiven_seq min_readable_seq查询同样有效读取seq given_seqmin_readable_seq是存储引擎当前可证明支持的最小可读序列下界。在实现中它就是 Region 已刷新的序列前沿version.flushed_sequence。4.2 源码实现validate_sequence_fences存储引擎侧的实现位于 src/mito2/src/engine.rs 的validate_sequence_fences对memtable_min_sequence下界以flushed_sequence为min_readable_seq若given_seq min_readable_seq则触发IncrementalQueryStaleSnafu { region_id, given_seq, min_readable_seq }src/mito2/src/engine.rs对memtable_max_sequence上界当快照上界 H 已被刷入 SST 时存储引擎无法对 SST 扫描应用 memtable 级序列上界因此若given_seq flushed_sequence会触发SnapshotFenceStale错误由 Flow 重新绑定带围栏的修复扫描避免读到 H 之后的行源码还支持exact_sequence_range模式要求行级精确 delta 时若能力不可用则宁可报SequenceRangeUnsupported也不静默返回近似超集。4.3 错误码与测试覆盖IncrementalQueryStale被映射为StatusCode::RequestOutdatedsrc/mito2/src/error.rs客户端可据此判定游标过期。存储引擎侧的扫描测试覆盖了三分支语义见 src/mito2/src/engine/scan_test.rs、src/mito2/src/engine/scan_test.rs 与 src/mito2/src/engine/scan_test.rs 等用例。4.4 一个关键澄清flush 边界不是独立的回退类别RFC 特别澄清比检查点更新的行已经越过 memtable→SST 冲刷边界、无法再证明序列精确排他这一情形在 v1 中并不是独立的回退类别而是增量游标变得 stale 的一种具体方式。这一点很微妙但很重要它把数据落盘与游标过期统一到同一个IncrementalQueryStale错误路径上从而大幅简化了回退状态机的设计。五、核心设计四水位线返回与检查点推进5.1 终端 metrics 中的分区域水位线增量查询的成功结果需要携带可推进检查点的证据RFC 扩展查询 metrics增加可选的分区域水位线映射region_latest_sequences: Vec(region_id: u64, latest_sequence: u64)在实现中这一信息以region_watermarks字段RegionWatermarkEntry在查询结果/客户端间传递见 src/flow/src/batching_mode/frontend_client.rs 与 src/flow/src/batching_mode/task/test.rs 的构造与断言。5.2 三条硬性规则RFC 对水位线的使用规定了三条规则全部服务于正确性只有成功查询的终端 metrics 才能推进检查点——失败或中断的轮次不得推进多 Region 查询中水位线必须是完整映射或干脆缺席——不允许只覆盖部分 Region 的半吊子水位线用于推进当正确性无法证明时业务行可以返回但水位线必须缺席——即业务数据照常产出但系统不会假装有了可推进检查点的证据。5.3 状态机中的水位线推进实现Flow 任务状态src/flow/src/batching_mode/state.rs中CheckpointMode枚举只有两个状态FullSnapshot与Incrementalsrc/flow/src/batching_mode/state.rsadvance_checkpoints(watermark_map)用完整的全量查询水位线替换检查点映射并从 FullSnapshot 切换到 Incrementalsrc/flow/src/batching_mode/state.rsadvance_incremental_checkpoints_with_participation(participating_regions, watermark_map)只推进本轮参与 Region 的检查点适用于增量 delta 查询src/flow/src/batching_mode/state.rsmark_full_snapshot()无条件回到 FullSnapshot 模式src/flow/src/batching_mode/state.rs。六、核心设计五Flow 状态机、回退策略与持久化模型6.1 双状态机总览Flow 从全量模式启动按以下四步迁移全量查询成功且返回正确性水位线→ 进入增量模式checkpoint 建立增量查询成功且返回正确性水位线→ 推进检查点增量查询过期/失败→ 回退到全量模式全量查询未返回正确性水位线→ 停留在全量模式不切换。RFC 附带的 mermaid 状态图完整呈现了这一模型6.2 确定性回退策略Fallback Policy回退到全量模式是确定性的由以下任一条件触发收到IncrementalQueryStale增量查询因执行错误而失败增量查询成功但参与 Region 的水位线缺席或不完整。回退后的行为策略同样确定失败/不完整的轮次不推进任何检查点下一轮对受影响的 Flow/时间窗口切换到全量模式只有当一次全量查询成功且返回完整正确性水位线映射后才允许回到增量模式。实现中回退原因被枚举为FlowQueryFallbackReasonsrc/flow/src/batching_mode/checkpoint.rs包含MissingRegionWatermark、IncompleteRegionWatermark、DirtyBacklogPending、StaleCursor、SnapshotFenceExpired、IncrementalQueryFailure、QueryFailure、IncrementalDisabled等类别并配有as_label()便于指标打点每轮查询结果会被归约为FlowCheckpointDecisionAdvancedFromFullSnapshot/AdvancedIncremental/ContinuedFencedRepair/ 回退等用于驱动下一轮执行。一个额外的实现细节值得注意disable_incremental()会永久禁用该 Flow 的增量模式例如当查询形状不是可增量重写的聚合查询时并立即回退全量——这意味着某些无法增量化的 Flow 一经判定就长期保持全量模式避免反复尝试。6.3 持久化与恢复模型Persistence and Recoveryv1 刻意保持进度游标轻量采取正确性优先的持久化取舍水位线/检查点只保存在 flownode 内存中v1 不单独持久化冷启动或 flownode 重启后Flow 总是先通过一次全量快照读重新建立进度只有该轮全量查询成功并返回完整正确性水位线映射后才恢复增量模式序列精确的增量正确性仅限于仍可见于 memtable 的行一旦相关行被刷入 SST系统无法仅凭seq checkpoint证明精确排他SST 缺行级序列元数据该情形下正确行为是回退全量重算而不是继续尽力而为的增量扫描。这个模型的本质是用重启即全量重算一次的代价换取状态机在任意时刻都只依赖可证明的正确性从而彻底避免持久化检查点与内存状态之间的一致性陷阱。七、分布式与兼容性要求RFC 明确列出分布式路径的三大要求分布式路径必须端到端保持分 Region 的快照/读边界语义snapshot_seqs与flow.*选项必须同时正确传递——其中snapshot_seqs指分 Region 快照上界映射region_id - sequence它定义了本轮查询最多看到什么新增 metrics 字段必须向后兼容——旧客户端应能忽略未知字段。这一设计对应了仓库中从 src/flow/src/batching_mode/frontend_client.rs客户端侧的扩展与水位线解析到 src/query/src/options.rsFlowQueryExtensions解析再到 src/mito2/src/engine.rs存储引擎围栏校验的完整链路扩展键在查询上下文中被解析与校验随分布式查询传递最终在存储引擎侧施加序列围栏再把水位线随结果带回 Flow。八、落地路线Rollout Plan与风险8.1 Phase 1MVP正确性优先新增扩展常量与解析新增增量扫描映射与 stale 检测新增水位线 metrics 字段与终端水位线推进检查点路径完成 standalone 与 distributed 的贯通passthrough。8.2 Phase 2性能与可观测性结合序列/水位线上下文改进批处理键batching key策略优化水位线序列化开销新增指标增量命中率incremental hit rate、回退率fallback rate、回退窗口大小fallback window size。8.3 测试计划RFC 规划了五类测试仓库现状已经覆盖了大部分增量边界与 stale 检测的单测——见 src/query/src/options.rs扩展解析与 src/mito2/src/engine/scan_test.rsstale 三分支扩展映射与水位线语义的查询路径测试——见 src/query/src/dummy_catalog.rsscan request 从增量上下文映射Full→Incremental→Fallback 转换的 Flow 集成测试——见 src/flow/src/batching_mode/task/test.rs水位线构造与推进与 src/flow/src/batching_mode/state.rs状态机切换测试端到端快照/水位线传播的分布式测试新旧客户端-服务端组合的兼容性测试——src/flow/src/batching_mode/frontend_client.rs 已验证非法布尔值会报错属于参数校验类用例。8.4 风险清单RFC 坦承了四项主要风险理解它们有助于把握方案的适用边界边界语义不一致与混淆可能引发正确性缺陷——这也是 RFC 强调排他下界seq given_seq的原因分布式传播不完整会静默破坏水位线安全性Phase 2 优化落地前频繁回退可能降低吞吐memtable→SST 冲刷可能迫使比预期更多的全量重算直到 SST 具备更细粒度的序列跟踪能力。九、备选方案与最终结论9.1 被否决的备选方案RFC 明确比较了两个备选方案并说明拒绝理由把水位线放进业务行拒绝理由schema 污染——对应 Non-Goals 中不引入合成 watermark 列v1 就新增专用 Flight 消息类型拒绝理由为缩减范围而推迟。9.2 结论这份 RFC 为 Flow 批处理提供了一条务实且正确性优先的增量路径它复用了既有的序列扫描能力引入严格的 stale 处理并坚持只从可证明正确性的分区域水位线推进检查点。在仓库中该设计已从 RFC 演进为一整套可运行的实现——从查询扩展的解析校验src/query/src/options.rs、扫描决策src/query/src/dummy_catalog.rs、存储引擎围栏校验src/mito2/src/engine.rs到 Flow 任务的状态机与扩展构造src/flow/src/batching_mode/state.rs、src/flow/src/batching_mode/task/inc.rs与回退决策src/flow/src/batching_mode/checkpoint.rs形成了完整闭环。对希望深入 GreptimeDB Flow 增量计算机制的读者而言沿着这条链路阅读源码将能获得对序列、水位线、回退三者如何共同保障正确性的最直接理解。【免费下载链接】greptimedbThe open-source observability database. One columnar engine for metrics, logs, and traces, on object storage.项目地址: https://gitcode.com/GitHub_Trending/gr/greptimedb创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考