Aphasia

Cogito ergo sum.

引言:同一张表,不同的水流

赫拉克利特的河流残篇 B12 常被简化成“人不能两次踏入同一条河流”。但更接近残篇原意的表达是:踏入同一条河流的人们,遇到的是不断流来的不同水流。 Stanford Encyclopedia of Philosophy 对这一残篇的解释尤其值得注意:河流并不是在变化中失去自己;恰恰相反,它通过内容的持续更替保持为同一条河流。

数据库中的表也具有相似的双重性。我们把它称作同一张表,因为 Schema、主键、分区和业务语义保持着某种连续性;但表中的行不断被插入、删除和更新,承载它的文件、快照与存储位置也在持续变化。

Snapshot 回答的是“在某个时刻,这条河是什么样子”;CDC(Change Data Capture)回答的则是“从一个时刻走到另一个时刻,哪些水流经过了这里”。前者保存状态,后者描述状态如何成为现在。

这一区别看似直观,实现起来却并不只是把两个快照做一次集合相减。一个可重放、可消费的 Changelog 至少要回答四个问题:

  1. 身份:两条记录是否代表同一个逻辑实体?
  2. 顺序:多次变化以什么版本、序列号或事务顺序发生?
  3. 前像:更新和删除之前的值从哪里获得?
  4. 边界:Compaction、Vacuum 和 Retention 之后,变化还能否被恢复?

因此,CDC 的本质不是“找出两个状态的不同”,而是把状态之间的迁移编码为一组有身份、有顺序、有语义边界的事件。只要其中任何一项不稳定,UPDATE_BEFOREUPDATE_AFTERDELETE 就可能退化成猜测。

本文聚焦的是“如何从一张表派生 CDC”,而不是 Debezium、Flink CDC 等“如何把外部数据库的 WAL/Binlog 写入湖仓”。二者都处理变化,但一个从源日志捕获变化,另一个从表格式、快照和数据文件重建变化。下文将以公开文档和本地源码为依据,分析 Snowflake、BigQuery、Apache Iceberg、Delta Lake、Apache Hudi 与 Apache Paimon 的实现路线,并讨论这些经验如何落到托管式 MOR(Merge-on-Read)表格式的设计中。

版本说明:源码判断基于本文写作时检出的 Apache Iceberg 457854fb 与 Apache Paimon b782f24a;产品能力以 2026 年 8 月 25 日可访问的官方文档为准。文中明确区分“公开能力”“源码事实”和“方案推演”,避免把设计建议误写成现有产品能力。

1. 先说结论:CDC 的分歧在于何时物化变化

把“表如何产出 CDC”拆开,业界方案主要落在四个位置:

1写入事务                Compaction                 查询消费
2   │                         │                         │
3   ├─ 写时同步物化 ──────────┼────────────────────────►│
4   │                         ├─ Compaction 时物化 ────►│
5   │                         │                         ├─ 读时推导
6   └─ 异步旁路物化 ───────────────────────────────────►│
  • 读时推导(read-time derivation):不预先保存完整 Changelog,消费时根据版本、快照或隐藏元数据在线推导。
  • 写时同步物化(write-time materialization):在 Base 写入的同一事务中保存变更记录或补充日志。
  • Compaction 时物化(compaction-time materialization):等 Base 与 Delta 合并、新旧值同时可见时再产生 Changelog。
  • 异步旁路物化(async out-of-band materialization):用独立作业生成影子表或变更表。

四种路线没有绝对的优劣,它们只是把成本放在不同位置:

1写时物化:写放大 + 存储换低延迟读取与确定性
2读时推导:规划/读取放大换低写入成本与任意区间重放
3Compaction 物化:用可见性延迟换合并过程中的成本摊销
4异步物化:用系统解耦换版本对齐、延迟与重复存储

本文最重要的判断有五点:

  1. Before-image 才是 CDC 的成本中心。 INSERT 容易推导,UPDATE_BEFOREDELETE 的旧值决定了方案的复杂度。
  2. 稳定 Row Identity 是读时 CDC 的基础设施。 仅靠主键可以覆盖一部分场景,却难以处理无主键表、主键变更和同一区间内的多次更新。
  3. Compaction 不是纯物理维护。 一旦它删除了中间状态,就改变了未来能够回答的历史问题;若要语义透明,必须留下等价的变化证据。
  4. 最佳实现通常是混合路线。 按变化类型、可用旧值和新鲜度要求,把写时、读时与 Compaction 的长处组合起来。
  5. CDC 应被视为表的物理能力与保留契约,而不是附属导出任务。 Planner、Writer、Compactor、Cleaner 和 Reader 都必须理解同一套版本边界。

2. 分析框架:一条 Changelog 是怎样成立的

2.1 主轴:变化在什么时候被物化

评估 CDC 方案时,首先要问的是“变化行在哪个阶段形成”。越靠近写入端,越容易获得精确的操作语义和 Before-image,但也越容易增加写延迟与写放大;越靠近读取端,写入越轻,读取时需要还原的信息越多。

路线最适合的条件主要收益主要代价
读时推导元数据完整、版本可追溯、读侧可控零或低写放大,可按任意区间重放规划、回查与合并成本
写时同步物化Writer 同时持有新旧值强事务一致性,读取便宜写放大、存储和生命周期管理
Compaction 物化MOR 合并天然看见新旧状态摊销 Lookup/Merge 成本延迟依赖 Compaction,需保存锚点
异步旁路物化存量系统难以改造与主链路解耦,上线快异步滞后、版本对齐、重复计算与存储

2.2 真正难点:Before-image 从哪里来

一个完整的行级 Changelog 通常需要:

1+I  INSERT
2-U  UPDATE_BEFORE
3+U  UPDATE_AFTER
4-D  DELETE

+I+U 通常可以从新写入的数据获得,-U-D 却要求系统仍然能够访问旧值。主流实现的旧值来源可以归纳为:

  • Writer 在 Merge 时同时持有新旧记录,例如 Hudi 的 Merge Handle;
  • 读时沿 Position Delete、Row ID 或主键回查 Base;
  • Compaction 合并 Base 与 Delta 时进行 Lookup;
  • 稳定的 Row Lineage 在文件重写后仍延续逻辑行身份;
  • 上游直接输入完整 CDC,存储系统只负责保留而不负责重建。

这也是为什么“有主键”并不等于“天然有 CDC”。主键回答业务实体如何匹配,未必能表达一次物理删除属于更新、真正删除还是文件重写;无主键表更只能依赖稳定 Row ID、Position 或完整行比较。若一个区间内同一主键发生多次更新,只有版本顺序和稳定身份同时存在,系统才能恢复中间过程而不只是首尾差异。

2.3 不能忽略的其他维度

  • MOR 与 COW:MOR 通常保留 Delta/Delete,读时合并复杂;COW 读取简单,但文件重写会产生 Carryover Rows。
  • 事务与版本一致性:CDC 是否与 Base 在同一 Commit 中可见,还是需要跨任务对齐水位。
  • Full Delta 与 Min Delta:系统保存完整变化历史,还是只能从仍存活的文件推导最小差异。
  • 生命周期:Retention、Vacuum、Clean 和 Snapshot Expiration 如何约束可查询区间。
  • Schema Evolution:旧 Changelog 使用写入时 Schema,还是读取时的当前 Schema。
  • 成本位置:写放大、读放大、存储、随机回查、规划复杂度与延迟之间如何交换。
  • 消费语义:是一次性查询任意区间,还是维护一个会推进的 Offset;多个消费者是否彼此独立。

3. 主流系统逐一分析

3.1 Snowflake Streams:Offset 加隐藏变更元数据

Snowflake Streams 不保存一份完整表数据,而是保存 Source Table 的一个 Offset。基表启用 Change Tracking 后,系统通过 METADATA$ACTIONMETADATA$ISUPDATEMETADATA$ROW_ID 等隐藏列,在该 Offset 与当前版本之间返回变化;消费型 DML 提交后,Offset 向前推进。

这条路线的优点是:

  • 无独立全量副本,增加多个 Stream 的边际成本较低;
  • Offset 与表版本天然对齐;
  • 支持标准、Append-only 等不同消费语义。

它的代价也很明确:Stream 的可恢复性依赖 Base Table 的数据保留期。Offset 长期不推进并越过 Retention 后,Stream 会 Stale。换句话说,读时推导没有消灭历史存储成本,只是把它合并进 Base 的版本保留契约。

3.2 BigQuery CHANGES:托管存储上的时间区间函数

BigQuery CHANGES 通过 TVF 返回指定时间区间内的 INSERTUPDATEDELETE。表需要启用 enable_change_history,查询范围受 Time Travel Window 限制,并且单次查询区间存在约束。

BigQuery 和 Snowflake 共同说明了一点:当一个系统同时控制 Writer、元数据、存储和 Reader 时,读时变更推导可以成为产品级能力。但这里的“无独立 Changelog”并不等于“无成本”,隐藏历史、变更元数据与 Retention 仍然需要付费,只是这些成本不以用户可见的 CDC 文件出现。

3.3 Iceberg ChangelogScan:Snapshot 区间上的文件级变化

Apache Iceberg 的 IncrementalChangelogScan 把变化表达成 Snapshot 区间 (from, to] 上的扫描。本文检出的 BaseIncrementalChangelogScan 会:

  1. 取两个 Snapshot 之间的祖先链;
  2. 跳过 REPLACE 类型 Snapshot,避免把纯 Compaction 误报为业务变化;
  3. 根据 Manifest Entry 的 ADDED / DELETED 状态生成 Added Rows 或 Deleted Data File 任务。

但源码边界比“支持 Iceberg CDC”这句概括严格得多:只要区间内存在 Delete Manifest,当前实现就抛出 UnsupportedOperationException。因此这条 Core 路径仍然主要面向 COW,并没有把 Equality Delete 或 Position Delete 自动翻译为完整 MOR 行级 CDC。

COW 还有 Carryover 问题:更新一行可能重写整个文件,文件级 -D/+I 会把未变化的行也带入 Changelog。Spark 侧的 ChangelogIterator 负责 removeCarryovers,并可按 Identifier Columns 把 Delete/Insert 配对为 Update。也就是说:

1Iceberg Core:发现文件级变化
2Spark Engine:消除 Carryover,补充行级 Net Change 语义

这一区分很关键。表格式提供了可追溯的物理事实,执行引擎仍需完成一部分语义恢复。

3.4 Iceberg v3 Row Lineage:让逻辑行穿过文件重写

Iceberg Format v3 的 Row Lineage 引入 _row_id_last_updated_sequence_number:新行获得稳定 Row ID,更新后的行继承原有身份并更新序列号。通过提交阶段的继承规则,Writer 不必提前知道最终 Snapshot 信息,也不必为赋值重写文件。

Row Lineage 解决的是比主键更底层的问题:

  • 无主键表仍可识别同一逻辑行;
  • 文件重写后,行身份不随 Position 改变;
  • 更新配对不再完全依赖 Identifier Columns 的启发式匹配;
  • 变化可以用 (row_id, last_updated_sequence_number) 建立稳定顺序。

不过,Row Lineage 与当前 Core ChangelogScan 仍是需要进一步接通的两套能力,不能把规范中的 Lineage 直接等同于已经完成的 MOR CDC 实现。它更像是下一代读时 Changelog 的必要地基。

3.5 Iceberg 社区讨论:Lineage 解决身份,但没有凭空生成历史

Iceberg 社区围绕 Row Lineage 的讨论经历了一个重要变化。早期提案把它设计成 v3 的可选能力;在 [DISCUSS] Row lineage required for v3 线程中,社区倾向于让所有 v3 表都携带 Lineage,以避免 Reader 必须猜测某个 Writer 是否遵守同一套行身份规则。讨论同时保留了一个现实边界:Writer 可以把一次变化建模为“保留 Row ID 的修改”,也可以建模为“删除旧行并插入新行”。

这个边界最终体现在规范里:通过 Equality Delete 完成的更新无法继承旧 _row_id,因为这类 Writer 的目的就是避免读取旧数据。对 Lineage Reader 而言,它只能被解释为 DELETE + INSERT,不能无条件还原成 UPDATE_BEFORE + UPDATE_AFTER

社区实现也说明,Lineage 是一项全链路契约,而不是增加两个 Metadata Column 就结束了:

  • Spark 的 Preserve row lineage on compaction 和 Flink 的 Preserve row lineage in RewriteDataFiles 要求文件重写复制原有 _row_id_last_updated_sequence_number
  • 如果 Compaction 给未变化行重新分配 Row ID,纯物理重写会被 CDC 错判成全表删除再插入;
  • 当前仍在讨论中的 MOR ChangelogScan Delete 支持 需要按 Snapshot 维护 Existing Deletes 与 Added Deletes,避免一行被重叠的 Equality/Position Delete 重复输出;
  • 社区对 MOR CDC 的讨论还指出,Lineage Column 仍位于数据行中,Reader 通常还是要读取 Data File,并结合 Delete File 或 DV 才能找回已经消失的行。它改善配对的正确性,却不自动消除 I/O。

因此,基于 Lineage 的 CDC 应被拆成三个独立问题:

1Change Discovery:哪些文件、Delete 或 DV 在区间内发生变化?
2Row Identity:变化前后的记录是不是同一个逻辑行?
3Event Reconstruction:它应输出 INSERT、DELETE,还是 UPDATE 前后像?

Iceberg Snapshot、Manifest 与 Delete/DV 主要回答第一个问题,Row Lineage 回答第二个问题,CDC Reader 才负责第三个问题。缺少其中任何一层,都无法形成完整方案。

3.6 Snowflake 路线:把 Iceberg v3 Lineage 接入 Streams

Snowflake 提供了当前最完整的产品化参考。官方 Iceberg v3 说明明确表示,Dynamic Iceberg Tables 与 Streams 会利用 Row Lineage 在开放表上完成增量处理,并覆盖 INSERTUPDATEDELETEMERGERow Lineage 使用约束进一步要求外部 Writer 遵守 Iceberg v3 规则:新行的 _row_id 应留空并由继承机制赋值,COW 复制的旧行必须保留 _row_id

Snowflake 没有公开其内部 Reader 的全部实现,因而不能把下面的细节当成源码事实;但它公开的行为足以说明这条路线的核心契约:

  1. Stream 保存消费 Offset,而不是保存一份全量 CDC 表;
  2. Iceberg v3 提供跨 Snapshot 稳定的 Row Identity;
  3. Snapshot/Manifest 和 Delete/DV 界定候选变化范围;
  4. Reader 使用 Row ID 配对前后状态,过滤只发生文件搬迁的 Carryover Rows;
  5. Retention 决定 Offset 之后的历史是否仍可重放。

这条路线比“对两个 Snapshot 做主键 Diff”更可靠,因为身份属于表格式,而不是由某个 Reader 临时选择的业务列;它也比写时保存完整 CDC 更轻,因为 Compaction 只需保留 Lineage,不必把未变化行误写成变化事件。

3.7 Delta CDF:从写时物化走向 Row Identity 驱动的双轨路线

传统 Delta Change Data Feed 是典型的混合方案。启用 delta.enableChangeDataFeed 后,UPDATEDELETEMERGE 可在同一事务中写入带 _change_type 等元数据的变更数据;能够从 Transaction Log 直接推导的操作则无需重复保存。消费者通过 table_changes 读取版本区间。

这条 Legacy CDF 路线的优点是 Base 与 CDC 同事务提交、读侧便宜;代价是更新路径的写放大,并且 CDC 的可用区间仍受 Vacuum 与历史保留控制。

值得关注的是,Databricks 当前又提供了 Automatic Change Data Feed:在满足运行时与 Unity Catalog 等条件时,读取端可以利用 Delta Row Tracking 或 Iceberg v3 Row Lineage 动态生成变化,而不要求每张表预先开启 Legacy CDF。这个能力仍有 Public Preview 与专有 Reader 等边界,但它揭示了清晰的演化方向:

1没有稳定行身份
2    └─► 更依赖写时保存完整变化
3
4具备 Row Tracking / Row Lineage
5    └─► 可以把更多 CDC 成本后移到读取端

这比简单地把 Delta 归类成“写时物化”更准确:同一个系统可以同时提供 Legacy 写时 CDF 与基于稳定 Row Identity 的 Automatic 读时 CDF,并让用户按兼容性、成本和延迟选择。

3.8 Hudi CDC:在 Merge 时捕获新旧值

Apache Hudi 的 CDC 查询RFC-51 代表写时/合并时物化路线。Hudi 可以通过 Supplemental Logging 控制 CDC 记录包含多少信息:

  • KEY_OP:只保存 Key 和操作类型,读取时补数据;
  • DATA_BEFORE:额外保存 Before-image;
  • DATA_BEFORE_AFTER:同时保存 Before-image 与 After-image。

设计的核心并不是“总要写一份完整 CDC”,而是根据读取成本选择最小必要日志。生成 Base File 的 Merge 阶段能够同时观察旧记录与新记录,因此适合保存准确前后像;对 MOR 中尚未合并的新日志,也可以结合已有 Log Block 推导变化。Hudi 因而更接近“持久化 CDC + 在线推导”的混合实现,而不是一个单一的同步写文件模型。

它给工程设计的最大提醒是生命周期:CDC Block、Log File、Compaction 和 Cleaner 必须共享一致的时间边界。若物化的变化证据早于消费 Offset 被清理,再准确的 CDC 生成逻辑也无法完成重放。

3.9 Paimon Changelog Producer:把成本位置显式交给表属性

Apache Paimon Changelog Producer 把策略直接暴露为表属性。当前主要包含四种模式:

  • none:默认模式,不额外生成完整 Changelog;
  • input:保留上游输入的完整 Changelog,适合 Flink CDC 等来源;
  • lookup:在 Compaction 中 Lookup 旧值,产生完整前后像;
  • full-compaction:在 Full Compaction 时比较前后结果,延迟更高但行为集中。

需要特别澄清:first-row 是一种 Merge Engine,不是第五种 changelog-producer。它可以与 Changelog 行为共同讨论,但不能放在同一个枚举层次。

Paimon 的价值在于把 Before-image 的成本位置说得非常直白:上游能提供就使用 input,愿意支付 Lookup Compaction 就使用 lookup,可以容忍更高延迟就使用 full-compaction。代价是配置复杂度和 Compaction 性能开销;官方文档也明确提示,启用 Changelog Producer 可能降低 Compaction 性能。

3.10 源侧 CDC:与表内 Changelog 分清边界

Debezium、Flink CDC 读取的是数据库 WAL、Redo Log 或 Binlog,源系统已经提供操作顺序和事务信息。这类系统解决的是:

1源数据库日志 ──► 标准化变更事件 ──► 湖仓表

本文讨论的表内 CDC 则是:

1湖仓表的 Snapshot / File / Delete / Lineage ──► 变更事件

Paimon input 属于二者的连接点:表格式不重建 Before-image,而是保留上游已经给出的完整 Changelog。若上游只发送最新值,存储层之后仍然要面对旧值 Lookup、顺序和 Retention 问题。

4. 横向对比:没有“免费 CDC”,只有不同的成本归属

系统或能力主要生成时机Before-image 来源事务/版本一致性生命周期边界关键限制
Snowflake Streams读时推导隐藏变更元数据与历史状态表版本 OffsetBase RetentionOffset 过旧会 Stale
BigQuery CHANGES读时推导托管存储的变化历史时间区间Time Travel Window需启用 Change History,区间受限
Iceberg ChangelogScan读时推导Snapshot 与 Manifest 的文件变化Snapshot 区间Snapshot/File Retention当前 Core 主要支持 COW,MOR Delete 未接通
Iceberg v3 Row LineageCDC 身份层,不单独生成事件稳定 Row ID 与更新序列Snapshot Sequence依赖 Lineage 与旧文件共同保留Equality Delete 更新退化为 Delete + Insert,仍需读取变化文件
Snowflake Iceberg v3 Streams基于 Offset 的读时推导Iceberg Row Lineage 与表历史Stream OffsetBase Retention外部 Writer 必须正确维护 Lineage
Delta Legacy CDF写时为主、部分读时推导Writer 与 Transaction Log同一 Delta CommitVacuum/Retention更新路径存在写放大
Databricks Automatic CDF读时推导Delta Row Tracking / Iceberg Row LineageCommit 区间RetentionPreview、运行时和 Reader 边界
Hudi CDCMerge 时物化并结合读时推导Merge Handle、Log BlockHudi TimelineCleaner/CompactionCDC 与清理策略必须协同
Paimon输入或 Compaction 时物化上游 Changelog 或 LookupSnapshot/CompactionSnapshot/Changelog RetentionProducer 配置与 Compaction 成本

跨系统可以看到一条共同规律:云数仓更倾向在托管存储内部维护变化证据并读时推导;开放湖格式为了跨引擎确定性,往往提供可选的持久化 Changelog。随着 Row Tracking 和 Row Lineage 成熟,两者开始汇合——开放格式逐渐拥有稳定行身份,托管引擎也能在不牺牲语义的情况下把更多工作放到读取端。

5. MaxCompute Delta Table:公开能力与结构机会

MaxCompute Delta Table 提供 ACID、主键、增量查询、Time Travel 与 Compaction 等能力。公开文档描述的存储路径包含 Base 与 Delta 数据,更新和删除可以通过事务信息及行位置关联已有记录;增量查询 则面向指定时间区间读取增量结果。

这里必须划清事实边界:当前公开文档中的增量查询返回的是合并后的最新增量状态,并不等价于完整的 +I/-U/+U/-D CDC。 文档也说明完整更新状态的 CDC 仍属于后续能力。因此,下面不是对现有产品实现的描述,而是基于公开结构进行的设计推演。

5.1 结构机会一:删除指针可以成为 Before-image 的定位线索

Iceberg COW 的困难在于,一次单行更新会表现成旧文件删除与新文件增加,逻辑行配对在文件重写中被稀释。若 MOR Delta 中的删除记录能够通过事务与 Row Position 精确定位旧 Base Row,那么读取端就不必在整个旧文件中猜测哪一行被更新,而可以沿指针回查 Before-image。

不过,物理 Row Position 通常会在 Compaction 后变化。它适合定位“当前保留周期内的旧行”,却不能天然替代跨重写稳定的逻辑 Row ID。长期设计仍需明确:Compaction 后是保留映射、保存 Changelog 锚点,还是引入类似 Iceberg v3 的 Row Lineage。

5.2 结构机会二:MOR Merge 天然是新旧值相遇的时刻

MOR 读取或 Compaction 本来就要把 Base 与 Delta 合并、按主键处理更新和删除。这个过程同时拥有旧状态与新状态,与 Hudi Merge Handle、Paimon Lookup Compaction 所创造的条件相同。

因此,完整 CDC 不一定需要再启动一条独立任务重算差异。更自然的方式是让同一个 Merge Operator 同时产生两个输出:

1Base + Delta
234MOR Merge / Compaction
5    ├─► 最新表状态
6    └─► 有序 Changelog 或 Changelog Anchor

这不是“零成本”,但它避免了重复读取和重复匹配,把 CDC 成本摊销在本来就要发生的状态合并中。

5.3 异步影子表为什么只能是迁移路线

独立异步任务生成变更表具有改造小、系统解耦的优势,适合作为存量架构的过渡方案。但它长期会面对三个难题:

  1. Base 与 CDC 使用不同 Commit,消费者如何获得一致的版本水位;
  2. 任务失败或延迟时,如何证明没有丢失或重复变化;
  3. Base 的 Retention/Compaction 已推进而影子任务尚未消费时,旧值还能否恢复。

只要 CDC 是表的正式能力,这些问题最终仍会回到存储事务、版本元数据和生命周期管理中。旁路任务可以承载计算,却不应成为 CDC 正确性的唯一所有者。

6. 面向托管 MOR 表格式的四条候选路线

6.1 路线 A:以读时推导为主

对版本区间 (v0, v1] 扫描 Delta:新增行产生 +I,删除指针回查旧值产生 -D;同一稳定身份上的删除与新增按事务顺序组成 -U/+U

收益:

  • Base 写入路径几乎不增加负担;
  • 无独立 CDC 副本,任意区间可重放;
  • 版本与 Base 天然一致;
  • 适合低频 CDC 消费或较短增量区间。

挑战:

  • 大区间会积累规划、文件扫描与随机回查成本;
  • 同一主键多次更新必须保留事务顺序,不能只看首尾;
  • 无主键表需要稳定 Row ID;
  • Compaction 删除中间态后,跨边界的完整变化无法仅从最新 Base 恢复。

6.2 路线 B:按算子选择写时物化

Append 等操作继续读时推导;当 Writer 已经持有旧值时,在同一事务中保存必要的 Before/After 数据。这接近 Delta Legacy CDF 的经济性:不为容易推导的变化重复写数据,只为难以重建的变化付费。

收益: CDC 与 Base 强一致,消费延迟低,复杂更新读取便宜。

挑战: 批式 MOR Writer 未必会在写入时 Lookup 旧值;为 CDC 强制 Lookup 可能破坏延迟合并的收益。物化数据还必须与 Vacuum/Retention 同步治理。

6.3 路线 C:Compaction 生成锚点,读时推导补鲜

Compaction 负责把已经合并的历史段转换成完整 Changelog Anchor;Anchor 之后尚未合并的 Delta 则按路线 A 在线推导。读取区间可以拆成:

1(v0, v1]
2   = 已压实历史段:读取 Changelog Anchor
3   + 未压实新鲜段:在线扫描 Delta 并回查 Base

这条路线把完整性、新鲜度和成本分开处理:

  • Anchor 让 Compaction 不再抹去历史语义;
  • 在线段通常较短,控制读放大;
  • 新旧值匹配复用 Compaction 已有 Merge;
  • Anchor 与 Base 使用同一 Retention 契约,避免悬空引用。

主要代价是 Compactor 需要承担双输出、顺序与失败恢复逻辑。系统还必须保证 Base Commit 与 Anchor Commit 原子关联,否则只是把异步版本对齐问题换了一个位置。

6.4 路线 D:建立跨文件重写的 Row Lineage

稳定的 row_id 加更新序列号可以统一主键表与无主键表的变化身份。Compaction 继承 Row ID,真正的新行分配新 ID,更新推进 Last Updated Sequence。它不是单独的 CDC 生成路线,而是路线 A/C 的正确性增强层。

6.4.1 Writer 与 Commit 契约

参考 Iceberg v3 的继承机制,Writer 不直接竞争全局 Row ID,而是在 Commit 成功时由元数据层为新行分配不重叠的 ID 区间:

 1TableMetadata.next_row_id
 2 3 4Snapshot.first_row_id
 5 6 7Manifest.first_row_id
 8 910DataFile.first_row_id + row_position
111213Logical _row_id

新行把 _row_id_last_updated_sequence_number 留空,通过 File/Manifest/Snapshot 的继承信息在读取时补齐;更新行显式携带旧 _row_id,并把 Last Updated Sequence 留给当前 Commit 赋值;纯文件重写则同时复制两个字段,表示逻辑行没有改变。

为了支持多引擎写入,建议每个 Snapshot 额外记录可校验的 Lineage Capability,例如 Writer 是否完整保留 Row ID、是否把更新建模为 Delete/Insert。这与 Iceberg 邮件列表中“Reader 应能知道一个 Snapshot 是否正确保留 Row ID”的担忧一致。若某个区间包含不满足契约的 Snapshot,Reader 应显式降级为 DELETE + INSERT 或拒绝 Full CDC,不能静默猜测。

6.4.2 Snapshot-by-Snapshot 的 CDC Scan

CDC Reader 不能只比较区间两端。A → B → A 的首尾状态相同,但包含两次变化;正确计划需要按 Sequence Number 依次处理 (s0, s1] 中的每个业务 Snapshot:

1for snapshot in snapshots_between(s0, s1):
2    added_data    = newly_added_data_files(snapshot)
3    removed_data  = newly_removed_data_files(snapshot)
4    added_deletes = newly_added_delete_files_or_dvs(snapshot)
5
6    after_rows  = read_added_rows(added_data, added_deletes)
7    before_rows = read_affected_old_rows(removed_data, added_deletes)
8
9    emit classify_by_row_id(before_rows, after_rows, snapshot.sequence_number)

其中 REPLACE/Compaction Snapshot 不直接产生业务事件,但它们生成的文件仍可能是后续 Delete 的目标,所以 Planner 不能简单丢弃其物理状态;正确做法是忽略它的事件输出,同时让重写后的文件继续参与后续扫描。

6.4.3 按 Row ID 重建事件

对同一个 Commit,将 Before 与 After 按 _row_id 配对,可以得到稳定分类:

BeforeAfterLineage 关系CDC 输出
Row ID 在本 Snapshot 首次分配+I
旧 Row ID 被 Delete/DV 移除-D
Row ID 相同,Last Updated Sequence 推进-U+U
Row ID 和 Last Updated Sequence 均不变Carryover,不输出
业务主键相同但 Row ID 不同默认 -D+I;可选按主键计算 Net Update

最后一行正是 Equality Delete 的边界:如果 Writer 没有读取旧行,就无法把原 Row ID 交给新行。底层 CDC 必须诚实地输出 Delete/Insert;上层可以在用户明确提供 Identifier Columns 时再计算 Net Update,但那是 Engine 语义,不再是 Lineage 事实。

Before-image 也不会由 Row ID 凭空产生。对于 Position Delete/DV,Reader 要用 File Path + Position 定位旧行并投影 _row_id;对于被 COW 删除的 Data File,要应用此前已经生效的 Delete,避免输出早已删除的记录。然后才能按 Row ID 与新文件中的 After-image 做 Hash/Sort Merge。

6.4.4 输出顺序与语义粒度

建议每条事件至少携带以下系统列:

字段作用
_change_typeINSERTDELETEUPDATE_BEFOREUPDATE_AFTER
_commit_snapshot_id定位产生变化的 Snapshot
_commit_sequence_number建立跨 Snapshot 的全序
_change_ordinal保证同一 Commit 内 Before 排在 After 之前
_row_id配对同一逻辑行并支持幂等消费

Iceberg Snapshot 保存的是 Commit 后状态,因此通用 Reader 能保证的是 Commit 级 CDC。如果同一 Commit 内一行先更新为 B、再更新为 C,但中间态没有写入表格式,Reader 只能看到 A 到 C;Statement 级或操作日志级 CDC 必须由 Writer 额外物化,Lineage 无法恢复从未提交的 B。

6.4.5 规划与性能优化

Lineage 方案不应退化成每个 Snapshot 做一次全表 Join。可以分层裁剪:

  1. 用 Snapshot Summary 和 Manifest Entry 只选择有业务变化的 Snapshot;
  2. 用 Partition、Data Sequence Number 和 Delete/DV 引用裁剪候选文件;
  3. _last_updated_sequence_number 保存文件级 Lower/Upper Bounds,只读取与目标区间相交的 Row Group;
  4. 对 Position Delete/DV 优先按 File Path 定位,避免按全局主键回查;
  5. 在消费频繁时维护轻量的 Row ID → File/Row Group 索引,或把老区间压成 Changelog Anchor;
  6. 将 Before/After 按 Row ID 做桶内 Merge,避免全局 Shuffle。

Iceberg 社区对 MOR ChangelogScan 的讨论说明:重叠 Delete 去重、历史 Delete Index 和 Compaction Output 都会影响正确性与成本。Row Lineage 能让最后的配对更准确,却不能替代这些变化发现结构。因此工程目标不是“只扫 Lineage Column”,而是让 Manifest/Delete 完成候选裁剪,让 Row Lineage 完成精确归因。

6.4.6 Retention 与降级策略

一个 Snapshot 区间只有在以下证据同时存在时,才具备 Full CDC 能力:

1Snapshot/Manifest 连续
2    AND 旧 Data File 可读
3    AND Delete/DV 未被清理
4    AND Row Lineage 在重写中连续

若缺少旧文件,系统仍可能返回 After-only Change Feed;若 Lineage 中断,可以返回 Delete/Insert;若 Snapshot 链中断,则应报告 Offset 已过期。把这些状态明确成 Capability,比返回一份看起来完整、实际漏掉 Before-image 的结果更重要。

这项改造涉及 Writer、Commit、Compactor、Delete、Schema Projection 和 Reader 全链路,成本最高;但它把很多原本依赖启发式匹配的问题变成精确的等值关系,也能服务增量物化视图、行级审计和缓存失效等更广泛场景。

6.5 先定义事件契约,再选择物理路线

很多 CDC 设计从“能不能读到变化文件”开始,最后才发现上下游对事件语义的理解不同。更稳健的顺序是先定义消费契约:

契约维度必须回答的问题
Offset按 Snapshot、Commit Sequence、LSN 还是 Timestamp 定位?是否严格单调?
Ordering只保证单表顺序、分区顺序,还是事务内/跨表顺序?
IdentityRow ID、主键与文件位置哪个稳定?主键变更如何表达?
Image提供 After-only、Before/After,还是完整操作类型?
Schema事件按写入时 Schema 还是读取时 Schema 投影?不可兼容演化如何处理?
DeliveryAt-least-once 还是 Exactly-once?重复事件的稳定去重键是什么?
RetentionOffset 何时过期?过期后报错、降级还是重新 Bootstrap?

“Exactly-once CDC”尤其不能只由 Source 单方面承诺。Source 可以稳定重放 (table_id, commit_sequence, row_id, change_ordinal),Consumer 仍需把业务写入与 Offset 提交放在同一事务,或使用幂等 Upsert/去重表。否则 Source 不重复,Consumer 在写入成功、提交 Offset 前崩溃,恢复后仍会再次应用事件。

Delta Lake 的官方 CDF 文档也明确给出两个现实边界:CDF 只记录启用后的变化,并随表的 Retention/VACUUM 生命周期清理。这类边界必须进入 API 和监控,不能只放在运维手册里。

7. 推荐架构:把变化证据留在它最便宜出现的位置

对于拥有统一 Writer、Metadata、Compactor 与 Reader 的托管 MOR 系统,更均衡的方向是:

以路线 C 的 Compaction Changelog Anchor 保证历史完整性,以路线 A 的读时推导保证最新鲜区间;中期用路线 D 的 Row Lineage 统一行身份。

这套组合不是把三个机制简单堆叠,而是按信息出现的自然位置保存证据:

  1. Writer 已持有完整前后像时,可以选择写下最小必要日志;
  2. MOR Merge 首次同时看见新旧值时,生成稳定 Changelog;
  3. Compaction 即将删除中间态前,必须留下可重放的 Anchor;
  4. Reader 只推导 Anchor 之后的短区间;
  5. Cleaner 根据所有消费者水位与统一 Retention 回收 Base、Delta 和 Anchor。

一个可能的读取计划如下:

 1CDC Scan(from_version, to_version)
 2 3    ├─ Metadata Planner
 4    │    ├─ 定位可覆盖的 Changelog Anchors
 5    │    ├─ 定位尚未压实的 Delta Files
 6    │    └─ 校验 Snapshot / Retention 连续性
 7 8    ├─ Anchor Reader ────────────────┐
 9    │                                │
10    ├─ Delta Reader + Base Lookup ───┼─► 按 Row ID / PK 与 Sequence 合并
11    │                                │
12    └─ Schema Resolver ──────────────┘
131415                              +I / -U / +U / -D

7.1 正确性不变量

实现时至少应把以下条件写成可验证的不变量:

  • 对同一 Version Range 重读,结果确定且顺序稳定;
  • Base Commit 可见时,对应 Anchor/Delta 元数据也处于一致状态;
  • Compaction 前后的 CDC 结果在逻辑上等价;
  • Consumer Watermark 未越过 Retention 前,Cleaner 不得删除必需证据;
  • Schema Evolution 后仍能解释旧事件,不能静默错位;
  • 同一 Row ID 上的事件序列满足单调性;
  • 重试只产生幂等结果,不重复提交 Anchor。

7.2 需要用数据回答的性能问题

  • 不同区间长度下,Delta 文件数与 Planning Time 的关系;
  • Row ID/Position 回查是顺序读取、桶内 Lookup 还是随机 I/O;
  • Anchor 生成对 Compaction CPU、内存与写吞吐的放大比例;
  • 主键热点与同键多次更新对 Merge State 的影响;
  • Partial Update、聚合 Merge Engine 和 Schema Evolution 的额外成本;
  • 多少活跃消费者会把 Retention 水位拖长,增加多少存储。

策略可以由引擎根据这些指标自动选择,不必把 Paimon 式多个底层旋钮全部暴露给用户。用户更需要的是清晰的语义承诺,例如“可查询最近 N 天完整变化”和“端到端延迟不超过 M 分钟”。

8. 我的思考:CDC 是变化的物理代数

如果说 Snapshot 是名词,CDC 更像动词。名词描述“对象是什么”,动词描述“对象发生了什么”。两个 Snapshot 的差异只能告诉我们起点和终点不同;Changelog 还要说明中间经过了哪些动作,以及这些动作如何组合。

这使 CDC 很像一套变化的代数:

1State(v1) = Apply(State(v0), Changes(v0, v1])
2
3Changes(v0, v2]
4  = Changes(v0, v1] ⊕ Changes(v1, v2]

其中 Apply 不是普通集合运算, 也不是无条件拼接。它们成立需要稳定身份、版本顺序、操作语义和 Schema 解释。若系统只保存最终值,那么 0 → 1 → 0 在首尾快照上没有差异,却包含了两次真实变化;审计、增量聚合和下游触发器关心的正是这两次过程。

由此可以得到几个更一般的判断。

8.1 身份比差异更基础

没有稳定身份,“变化”就只是两批数据之间的相似性匹配。主键提供业务身份,Row ID 提供存储身份,Sequence Number 提供时间身份。三者并不互相替代:主键可能被修改,物理位置会被 Compaction 改写,而版本号本身又不能说明两条记录是否属于同一实体。

Iceberg v3 Row Lineage 和 Delta Row Tracking 之所以重要,不只是它们让 CDC 更快,而是它们把“这还是同一行”从执行引擎的推测提升为表格式的事实。

8.2 Compaction 必须证明自己语义透明

数据库通常把 Compaction 当成不改变查询结果的物理优化。对 Snapshot Query 而言,只要压实前后的最终状态一致,它就是透明的;对 CDC Query 而言,这个定义不够。如果中间更新被压掉,Compaction 已经改变了系统能够回答的问题集合。

因此,支持 CDC 的 Compaction 有两种选择:保留足够的原始历史,或者生成语义等价的 Changelog Anchor。任何删除变化证据的后台维护任务,都必须同时说明它保留了什么等价信息。 这条原则也适用于 Vacuum、Snapshot Expiration 和 Schema Rewrite。

8.3 Retention 是语义,不只是存储参数

用户看到的“保留 7 天”不应只是文件清理周期,而应是一项可验证的可重放承诺:在这 7 天内,任意合法版本区间能否返回完整 CDC?如果 Base 还在但 Delete File 已消失,或者 Changelog 在但引用的旧数据已清理,形式上的保留并没有带来语义上的可恢复。

真正的 Retention Watermark 应由 Snapshot、Delta、Delete、Lineage、Anchor 和所有消费者 Offset 的最小安全边界共同决定。

8.4 增量计算的核心是保存“可复用的过去”

增量计算并不是少算一点这么简单。它首先要知道过去的哪些结果仍然有效,哪些输入变化使它们失效,以及如何只重建受影响的部分。CDC 提供变化边界,物化状态提供可复用的过去,Lineage 则把二者连接起来。

这一思路也可以帮助理解 AI/LLM 系统中的缓存、增量索引和 Agent Memory:真正困难的往往不是缓存某个结果,而是确定身份、依赖关系和失效条件。没有稳定的变化语义,所谓“减少重复计算”很容易退化成复用过期状态。

所以我更愿意把 CDC 看成增量系统的基础语言。它不仅服务数据同步,还决定增量物化视图、流批一体、缓存失效、审计回放和在线特征更新能否共享同一套事实来源。

9. 结语

赫拉克利特的河流之所以仍是同一条河,不是因为水从未变化,而是因为变化遵循了一种可以辨认的连续性。表也一样:Snapshot 维持它作为“同一张表”的状态,CDC 则记录这种连续性内部发生的迁移。

从 Snowflake 和 BigQuery 的读时推导,到 Iceberg 的 Snapshot/Manifest 扫描与 Row Lineage,再到 Delta 的双轨 CDF、Hudi 的 Merge Logging 和 Paimon 的 Compaction Producer,各系统真正交换的都是同一组成本:何时获得 Before-image,在哪里保存身份和顺序,以及由谁承担历史保留。

对于托管 MOR 系统,最值得利用的不是另造一条孤立的影子链路,而是已有的 Merge、Compaction、事务元数据和删除指针。让 Compaction 在抹去中间态之前留下 Changelog Anchor,让 Reader 只推导最新鲜的一段,再用 Row Lineage 把逻辑身份贯穿文件重写,才能同时兼顾完整性、新鲜度与成本。

最终,一个成熟的 CDC 系统应该让下面的等式成为可验证的工程事实,而不只是一句设计愿望:

1旧状态 + 完整且有序的变化 = 新状态

参考资料与证据边界

本文对开源项目的实现判断来自上述源码与官方文档;对 MaxCompute 的现有能力仅采用公开文档,对 Changelog Anchor、读时回查与 Row Lineage 的组合均属于架构分析和方案推演,不代表产品已经实现或承诺上线。开源项目与云产品迭代较快,具体能力应以使用时版本为准。