Iceberg V3 流式 CDC 场景下 Equality Delete → Deletion Vector 实时转换机制深度调研
Apache Iceberg 的 Merge-on-Read 模式支持两种行级删除机制:
| 类型 | 原理 | 写入代价 | 读取代价 | 适用引擎 |
|---|---|---|---|---|
| Equality Delete | 按列值匹配(如 id=42)标记删除 | 极低(微秒级) | 高 — 需 join 所有候选 data files | Flink CDC |
| Position Delete / DV | 按物理位置 (file, row_position) | 中等(需先扫表定位) | 低 — O(1) bitmap 查找 | Spark |
核心约束:流式引擎不扫表。
Flink 收到 CDC 消息 DELETE id=3 时,只知道主键值,不知道该行存储在哪个 Parquet 文件的第几行。要获取物理位置需要全表扫描 — 对于 TB 级表,这会将毫秒级延迟变为分钟级,完全不可接受。
| 分支 | 内容 | 谁写 | 谁读 |
|---|---|---|---|
| staging 分支 | data files + equality deletes | IcebergSink (Writer) | ConvertEqualityDeletes task |
| target 分支 (main) | data files + DVs(无 equality delete) | ConvertEqualityDeletes task | 所有下游查询引擎 |
这种分支隔离确保 equality delete 永远不暴露给 reader,从根本上消除了读放大问题。
| 特性 | 描述 |
|---|---|
| V3 Only | 目标分支必须是 format-version=3(DV 是 V3 特性) |
| 分支隔离 | Writer 写 staging 分支,Reader 读 target 分支,equality delete 对 reader 不可见 |
| 幂等提交 | 通过 snapshot summary property 记录已处理的 staging snapshot id,crash 后不会重复提交 |
| 容错恢复 | Planner state 通过 aligned Flink checkpoint 持久化,committer 的 isAlreadyCommitted guard 保证幂等 |
| 互斥锁 | 通过 TableMaintenance 框架 coordination lock 与其他 maintenance task 互斥 |
| Reindex 感知 | 当 main 分支有外部 commit(如 Spark compaction)时自动触发重建索引 |
| 可配置 equality fields | ConvertEqualityDeletes.Builder.equalityFieldColumns(List<String>) |
| Phase-based 顺序保证 | 使用 watermark 保证 data indexing 在 delete resolution 之前完成 |
| 组件 | 指标 | 说明 |
|---|---|---|
| Planner | processedEqDeleteFileNum | 已处理的 equality delete 文件数 |
| Planner | processedStagingSnapshotNum | 已处理的 staging snapshot 数 |
| Planner | reindexCount | 重建索引次数 |
| Worker | indexedKeyNum | 索引中的 PK 数量 |
| Worker | resolvedDeleteNum | 已解析的 delete 数 |
| Committer | addedDvNum | 提交的 DV 文件数 |
| Committer | commitDurationMs | 提交耗时 |
MySQL/PostgreSQL → Flink CDC → IcebergSink (staging) → ConvertEqualityDeletes → target (DV only)
解决了"CDC 写入快但查询慢"的根本矛盾。下游 Spark/Trino/Snowflake 直接读取无 equality delete 的表。
电商库存、用户画像等频繁 UPDATE 场景,equality delete 堆积速度远超传统 compaction 频率。内联转换实时消化。
大批量行级删除(如"被遗忘权"请求),equality delete 文件数量爆炸。转换为 DV 后查询性能恢复正常。
Snowflake(外部 Iceberg 模式)和 Databricks 不支持 equality delete。转换后所有引擎可直接读取。
| 项目 | 值 |
|---|---|
| PR | #15996 |
| 作者 | Maximilian Michels (Iceberg Flink connector maintainer) |
| 创建日期 | 2026-04-16 |
| 合入分支 | apache:main(2026-07 合入) |
| 预计发布版本 | Iceberg 1.12.0(尚未发布,预计 2026 Q3-Q4) |
| Flink 版本 | Runtime Jar | 状态 |
|---|---|---|
| Flink 2.1 | iceberg-flink-runtime-2.1-1.12.0.jar | 预计支持 |
| Flink 2.0 | iceberg-flink-runtime-2.0-1.12.0.jar | 预计支持 |
| Flink 1.20 | iceberg-flink-runtime-1.20-1.12.0.jar | 预计支持 |
| Flink 1.19 及更早 | — | End of Life |
format-version=3。V2 表需要先升级到 V3:ALTER TABLE t SET TBLPROPERTIES ('format-version' = '3');
_row_id, _last_updated_sequence_number)| 方案 | 执行位置 | 延迟 | 需额外基础设施 | Reader 可见 eq-delete |
|---|---|---|---|---|
| ConvertEqualityDeletes (新) | Flink 作业内 | 实时(checkpoint 级) | 无 | 不可见 |
| Spark RewriteDataFiles | 外部 Spark 集群 | 分钟~小时(调度频率) | Spark 集群 | 窗口期可见 |
| S3 Tables 自动 Compaction | AWS 托管后台 | 分钟级(不完全可控) | 无(S3 Tables 内置) | 短暂可见 |
| IceStream (第三方) | 独立服务 | 接近实时 | 需部署运维 | 取决于配置 |
| Flink RewriteDataFiles task | Flink 作业内 | Compaction 触发频率 | 无 | Compaction 前可见 |
_row_id, _last_updated_sequence_number) 暂不保留| # | 链接 | 说明 |
|---|---|---|
| 1 | PR #15996 | ConvertEqualityDeletes 主 PR |
| 2 | 设计文档 | Google Doc 详细设计 |
| 3 | 邮件列表讨论 | Iceberg Dev 邮件列表 |
| 4 | Data Lakehouse Weekly 2026-07-18 | 社区报道 |
| 5 | The Equality Delete Problem | 问题深度分析 (RisingWave) |
| 6 | Flink TableMaintenance 文档 | Iceberg 官方文档 |
| 7 | Issue #11122 | V3 Position Delete 改进 Umbrella Issue |
| 8 | IceStream | 第三方 eq-to-DV 转换方案 |
| 9 | Streaming Updates in Iceberg | Position Deletes at Scale (Etleap) |
| 10 | Multi-Engine Support | Iceberg 多引擎版本兼容矩阵 |
调研时间:2026-07-21 | 作者:Xiao Huang | 基于 PR #15996 及社区公开讨论整理