引言:同一张表,不同的水流
赫拉克利特的河流残篇 B12 常被简化成“人不能两次踏入同一条河流”。但更接近残篇原意的表达是:踏入同一条河流的人们,遇到的是不断流来的不同水流。 Stanford Encyclopedia of Philosophy 对这一残篇的解释尤其值得注意:河流并不是在变化中失去自己;恰恰相反,它通过内容的持续更替保持为同一条河流。
数据库中的表也具有相似的双重性。我们把它称作同一张表,因为 Schema、主键、分区和业务语义保持着某种连续性;但表中的行不断被插入、删除和更新,承载它的文件、快照与存储位置也在持续变化。
Snapshot 回答的是“在某个时刻,这条河是什么样子”;CDC(Change Data Capture)回答的则是“从一个时刻走到另一个时刻,哪些水流经过了这里”。前者保存状态,后者描述状态如何成为现在。
这一区别看似直观,实现起来却并不只是把两个快照做一次集合相减。一个可重放、可消费的 Changelog 至少要回答四个问题:
- 身份:两条记录是否代表同一个逻辑实体?
- 顺序:多次变化以什么版本、序列号或事务顺序发生?
- 前像:更新和删除之前的值从哪里获得?
- 边界:Compaction、Vacuum 和 Retention 之后,变化还能否被恢复?
因此,CDC 的本质不是“找出两个状态的不同”,而是把状态之间的迁移编码为一组有身份、有顺序、有语义边界的事件。只要其中任何一项不稳定,UPDATE_BEFORE、UPDATE_AFTER 和 DELETE 就可能退化成猜测。
本文聚焦的是“如何从一张表派生 CDC”,而不是 Debezium、Flink CDC 等“如何把外部数据库的 WAL/Binlog 写入湖仓”。二者都处理变化,但一个从源日志捕获变化,另一个从表格式、快照和数据文件重建变化。下文将以公开文档和本地源码为依据,分析 Snowflake、BigQuery、Apache Iceberg、Delta Lake、Apache Hudi 与 Apache Paimon 的实现路线,并讨论这些经验如何落到托管式 MOR(Merge-on-Read)表格式的设计中。
版本说明:源码判断基于本文写作时检出的 Apache Iceberg
457854fb与 Apache Paimonb782f24a;产品能力以 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异步物化:用系统解耦换版本对齐、延迟与重复存储
本文最重要的判断有五点:
- Before-image 才是 CDC 的成本中心。
INSERT容易推导,UPDATE_BEFORE和DELETE的旧值决定了方案的复杂度。 - 稳定 Row Identity 是读时 CDC 的基础设施。 仅靠主键可以覆盖一部分场景,却难以处理无主键表、主键变更和同一区间内的多次更新。
- Compaction 不是纯物理维护。 一旦它删除了中间状态,就改变了未来能够回答的历史问题;若要语义透明,必须留下等价的变化证据。
- 最佳实现通常是混合路线。 按变化类型、可用旧值和新鲜度要求,把写时、读时与 Compaction 的长处组合起来。
- 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$ACTION、METADATA$ISUPDATE 与 METADATA$ROW_ID 等隐藏列,在该 Offset 与当前版本之间返回变化;消费型 DML 提交后,Offset 向前推进。
这条路线的优点是:
- 无独立全量副本,增加多个 Stream 的边际成本较低;
- Offset 与表版本天然对齐;
- 支持标准、Append-only 等不同消费语义。
它的代价也很明确:Stream 的可恢复性依赖 Base Table 的数据保留期。Offset 长期不推进并越过 Retention 后,Stream 会 Stale。换句话说,读时推导没有消灭历史存储成本,只是把它合并进 Base 的版本保留契约。
3.2 BigQuery CHANGES:托管存储上的时间区间函数
BigQuery CHANGES 通过 TVF 返回指定时间区间内的 INSERT、UPDATE 与 DELETE。表需要启用 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 会:
- 取两个 Snapshot 之间的祖先链;
- 跳过
REPLACE类型 Snapshot,避免把纯 Compaction 误报为业务变化; - 根据 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 在开放表上完成增量处理,并覆盖 INSERT、UPDATE、DELETE 与 MERGE。Row Lineage 使用约束进一步要求外部 Writer 遵守 Iceberg v3 规则:新行的 _row_id 应留空并由继承机制赋值,COW 复制的旧行必须保留 _row_id。
Snowflake 没有公开其内部 Reader 的全部实现,因而不能把下面的细节当成源码事实;但它公开的行为足以说明这条路线的核心契约:
- Stream 保存消费 Offset,而不是保存一份全量 CDC 表;
- Iceberg v3 提供跨 Snapshot 稳定的 Row Identity;
- Snapshot/Manifest 和 Delete/DV 界定候选变化范围;
- Reader 使用 Row ID 配对前后状态,过滤只发生文件搬迁的 Carryover Rows;
- Retention 决定 Offset 之后的历史是否仍可重放。
这条路线比“对两个 Snapshot 做主键 Diff”更可靠,因为身份属于表格式,而不是由某个 Reader 临时选择的业务列;它也比写时保存完整 CDC 更轻,因为 Compaction 只需保留 Lineage,不必把未变化行误写成变化事件。
3.7 Delta CDF:从写时物化走向 Row Identity 驱动的双轨路线
传统 Delta Change Data Feed 是典型的混合方案。启用 delta.enableChangeDataFeed 后,UPDATE、DELETE 和 MERGE 可在同一事务中写入带 _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 | 读时推导 | 隐藏变更元数据与历史状态 | 表版本 Offset | Base Retention | Offset 过旧会 Stale |
BigQuery CHANGES | 读时推导 | 托管存储的变化历史 | 时间区间 | Time Travel Window | 需启用 Change History,区间受限 |
| Iceberg ChangelogScan | 读时推导 | Snapshot 与 Manifest 的文件变化 | Snapshot 区间 | Snapshot/File Retention | 当前 Core 主要支持 COW,MOR Delete 未接通 |
| Iceberg v3 Row Lineage | CDC 身份层,不单独生成事件 | 稳定 Row ID 与更新序列 | Snapshot Sequence | 依赖 Lineage 与旧文件共同保留 | Equality Delete 更新退化为 Delete + Insert,仍需读取变化文件 |
| Snowflake Iceberg v3 Streams | 基于 Offset 的读时推导 | Iceberg Row Lineage 与表历史 | Stream Offset | Base Retention | 外部 Writer 必须正确维护 Lineage |
| Delta Legacy CDF | 写时为主、部分读时推导 | Writer 与 Transaction Log | 同一 Delta Commit | Vacuum/Retention | 更新路径存在写放大 |
| Databricks Automatic CDF | 读时推导 | Delta Row Tracking / Iceberg Row Lineage | Commit 区间 | Retention | Preview、运行时和 Reader 边界 |
| Hudi CDC | Merge 时物化并结合读时推导 | Merge Handle、Log Block | Hudi Timeline | Cleaner/Compaction | CDC 与清理策略必须协同 |
| Paimon | 输入或 Compaction 时物化 | 上游 Changelog 或 Lookup | Snapshot/Compaction | Snapshot/Changelog Retention | Producer 配置与 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
2 │
3 ▼
4MOR Merge / Compaction
5 ├─► 最新表状态
6 └─► 有序 Changelog 或 Changelog Anchor
这不是“零成本”,但它避免了重复读取和重复匹配,把 CDC 成本摊销在本来就要发生的状态合并中。
5.3 异步影子表为什么只能是迁移路线
独立异步任务生成变更表具有改造小、系统解耦的优势,适合作为存量架构的过渡方案。但它长期会面对三个难题:
- Base 与 CDC 使用不同 Commit,消费者如何获得一致的版本水位;
- 任务失败或延迟时,如何证明没有丢失或重复变化;
- 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 │
9 ▼
10DataFile.first_row_id + row_position
11 │
12 ▼
13Logical _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 配对,可以得到稳定分类:
| Before | After | Lineage 关系 | 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_type | INSERT、DELETE、UPDATE_BEFORE、UPDATE_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。可以分层裁剪:
- 用 Snapshot Summary 和 Manifest Entry 只选择有业务变化的 Snapshot;
- 用 Partition、Data Sequence Number 和 Delete/DV 引用裁剪候选文件;
- 为
_last_updated_sequence_number保存文件级 Lower/Upper Bounds,只读取与目标区间相交的 Row Group; - 对 Position Delete/DV 优先按 File Path 定位,避免按全局主键回查;
- 在消费频繁时维护轻量的 Row ID → File/Row Group 索引,或把老区间压成 Changelog Anchor;
- 将 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 | 只保证单表顺序、分区顺序,还是事务内/跨表顺序? |
| Identity | Row ID、主键与文件位置哪个稳定?主键变更如何表达? |
| Image | 提供 After-only、Before/After,还是完整操作类型? |
| Schema | 事件按写入时 Schema 还是读取时 Schema 投影?不可兼容演化如何处理? |
| Delivery | At-least-once 还是 Exactly-once?重复事件的稳定去重键是什么? |
| Retention | Offset 何时过期?过期后报错、降级还是重新 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 统一行身份。
这套组合不是把三个机制简单堆叠,而是按信息出现的自然位置保存证据:
- Writer 已持有完整前后像时,可以选择写下最小必要日志;
- MOR Merge 首次同时看见新旧值时,生成稳定 Changelog;
- Compaction 即将删除中间态前,必须留下可重放的 Anchor;
- Reader 只推导 Anchor 之后的短区间;
- 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 ──────────────┘
13 │
14 ▼
15 +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旧状态 + 完整且有序的变化 = 新状态
参考资料与证据边界
- Heraclitus — Stanford Encyclopedia of Philosophy
- Snowflake — Introduction to Streams
- BigQuery —
CHANGESTime Series Function - Apache Iceberg — IncrementalChangelogScan Javadoc
- Apache Iceberg — BaseIncrementalChangelogScan Source
- Apache Iceberg — ChangelogIterator Javadoc
- Apache Iceberg — Format Specification
- Apache Iceberg Dev List — Row Lineage Required for v3
- Apache Iceberg Community — Preserve Row Lineage on Spark Compaction
- Apache Iceberg Community — Preserve Row Lineage in Flink RewriteDataFiles
- Apache Iceberg Community — MOR ChangelogScan Delete Support
- Snowflake — Apache Iceberg v3 Support
- Snowflake — Use Row Lineage with Iceberg Tables
- Databricks — Use Change Data Feed
- Databricks SQL —
table_changesFunction - Apache Hudi — RFC-51: Change Data Capture
- Apache Hudi — Change Data Capture Queries
- Apache Paimon — Changelog Producer
- MaxCompute — Delta Tables
- MaxCompute — Incremental Query
- MaxCompute — Update and Delete
本文对开源项目的实现判断来自上述源码与官方文档;对 MaxCompute 的现有能力仅采用公开文档,对 Changelog Anchor、读时回查与 Row Lineage 的组合均属于架构分析和方案推演,不代表产品已经实现或承诺上线。开源项目与云产品迭代较快,具体能力应以使用时版本为准。