Flink ConvertEqualityDeletes

Iceberg V3 流式 CDC 场景下 Equality Delete → Deletion Vector 实时转换机制深度调研

Iceberg 1.12.0 (未发布) PR #15996 已合入 main Flink 1.20 / 2.0 / 2.1
目录

1. 背景与问题

Equality Delete 的本质矛盾

Apache Iceberg 的 Merge-on-Read 模式支持两种行级删除机制:

类型原理写入代价读取代价适用引擎
Equality Delete按列值匹配(如 id=42)标记删除极低(微秒级) — 需 join 所有候选 data filesFlink CDC
Position Delete / DV按物理位置 (file, row_position)中等(需先扫表定位) — O(1) bitmap 查找Spark

为什么 Flink 只能写 Equality Delete?

核心约束:流式引擎不扫表。

Flink 收到 CDC 消息 DELETE id=3 时,只知道主键值,不知道该行存储在哪个 Parquet 文件的第几行。要获取物理位置需要全表扫描 — 对于 TB 级表,这会将毫秒级延迟变为分钟级,完全不可接受。

Equality Delete 带来的问题

2. ConvertEqualityDeletes 方案设计

核心思路:在 Flink 作业内使用 RocksDB State 维护主键索引,流式将 equality delete 解析为 position,直接以 Deletion Vector (Puffin) 格式输出。下游 Reader 永远看不到 equality delete。

分支隔离策略

分支内容谁写谁读
staging 分支data files + equality deletesIcebergSink (Writer)ConvertEqualityDeletes task
target 分支 (main)data files + DVs(无 equality delete)ConvertEqualityDeletes task所有下游查询引擎

这种分支隔离确保 equality delete 永远不暴露给 reader,从根本上消除了读放大问题。

3. 详细工作流

IcebergSink (Writer) │ ▼ commit to staging branch ┌─────────────────────────────────────────────────────────┐ │ ConvertEqualityDeletes Pipeline │ │ │ │ ┌─────────────────┐ │ │ │ 1. Planner (p=1)│ 扫描 staging 最老未处理 snapshot │ │ │ 发出 ReadCommands: │ │ │ Phase 0: main data files → 构建 PK index │ │ │ Phase 1: staging eq-delete files → 解析 │ │ │ Phase 2: staging data files → 更新 index │ │ └────────┬────────┘ │ │ ▼ │ │ ┌─────────────────┐ │ │ │ 2. Reader (p=N) │ 读取文件内容,发出 IndexCommands │ │ └────────┬────────┘ │ │ ▼ │ │ ┌─────────────────────────────┐ │ │ │ 3. Worker (p=N, keyed by PK)│ │ │ │ RocksDB State 维护 PK index shard │ │ │ eq-delete → 解析为 (filePath, position) │ │ └────────┬────────────────────┘ │ │ ▼ │ │ ┌──────────────────┐ │ │ │ 4. DVResolver (p=1)│ 按 data file 分组,合并已有 DV │ │ └────────┬─────────┘ │ │ ▼ │ │ ┌──────────────────┐ │ │ │ 5. DVMerger (p=N)│ 写入 DV (Puffin bitmap) │ │ └────────┬─────────┘ │ │ ▼ │ │ ┌────────────────────┐ │ │ │ 6. Committer (p=1) │ RowDelta commit 到 target 分支 │ │ │ data files + DVs (无 equality delete) │ │ └─────────────────────┘ │ └─────────────────────────────────────────────────────────┘ │ ▼ target branch (main) Query Engines (Spark / Trino / Snowflake / ...) → 只看到 data files + DVs ✅

4. 核心特性

V3
目标格式版本
6
Operator 阶段
RocksDB
PK Index 后端
Puffin
DV 文件格式
特性描述
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 fieldsConvertEqualityDeletes.Builder.equalityFieldColumns(List<String>)
Phase-based 顺序保证使用 watermark 保证 data indexing 在 delete resolution 之前完成

Metrics

组件指标说明
PlannerprocessedEqDeleteFileNum已处理的 equality delete 文件数
PlannerprocessedStagingSnapshotNum已处理的 staging snapshot 数
PlannerreindexCount重建索引次数
WorkerindexedKeyNum索引中的 PK 数量
WorkerresolvedDeleteNum已解析的 delete 数
CommitteraddedDvNum提交的 DV 文件数
CommittercommitDurationMs提交耗时

5. 使用场景

场景 1:Flink CDC 实时入湖(最核心)

MySQL/PostgreSQL → Flink CDC → IcebergSink (staging) → ConvertEqualityDeletes → target (DV only)

解决了"CDC 写入快但查询慢"的根本矛盾。下游 Spark/Trino/Snowflake 直接读取无 equality delete 的表。

场景 2:高频更新业务表

电商库存、用户画像等频繁 UPDATE 场景,equality delete 堆积速度远超传统 compaction 频率。内联转换实时消化。

场景 3:GDPR 合规删除

大批量行级删除(如"被遗忘权"请求),equality delete 文件数量爆炸。转换为 DV 后查询性能恢复正常。

场景 4:跨引擎数据湖

Snowflake(外部 Iceberg 模式)和 Databricks 不支持 equality delete。转换后所有引擎可直接读取。

6. 版本兼容与发布计划

项目
PR#15996
作者Maximilian Michels (Iceberg Flink connector maintainer)
创建日期2026-04-16
合入分支apache:main(2026-07 合入)
预计发布版本Iceberg 1.12.0(尚未发布,预计 2026 Q3-Q4)

Flink 版本兼容

Flink 版本Runtime Jar状态
Flink 2.1iceberg-flink-runtime-2.1-1.12.0.jar预计支持
Flink 2.0iceberg-flink-runtime-2.0-1.12.0.jar预计支持
Flink 1.20iceberg-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');

7. 社区演进时间线

2024-09
Issue #11122 — V3 Position Delete 改进总体设计提案(Umbrella Issue)
2024-10
Russell Spitzer 首次提出废弃 equality delete,因 Flink 无替代方案被搁置
2024-11
PR #11446 — Core 层为 DV 添加 content offset/size 到 DeleteFile(Flink metrics 同步适配)
2025-09
PR #14148 — Flink 支持读取 V3 row lineage 元数据字段 (_row_id, _last_updated_sequence_number)
2025-10
PR #14197 — Flink IcebergSink 支持写入 DV(position delete 路径)
2025-11
PR #14414 — DynamicIcebergSink 支持 DV 写入 + 多 format version 混合
2026-04-16
PR #15996 — ConvertEqualityDeletes maintenance task 创建
2026-05-19
Iceberg 1.11.0 发布 — V3 Production Ready,Flink TableMaintenance 框架完善
2026-07
PR #15996 合入 main,Maximilian Michels 确认跨 6 个 commit 完成
2026-07
huaxin gao 正式推动 V4 中废弃 equality delete(dev 邮件列表讨论)

8. 方案对比

方案执行位置延迟需额外基础设施Reader 可见 eq-delete
ConvertEqualityDeletes (新)Flink 作业内实时(checkpoint 级)不可见
Spark RewriteDataFiles外部 Spark 集群分钟~小时(调度频率)Spark 集群窗口期可见
S3 Tables 自动 CompactionAWS 托管后台分钟级(不完全可控)无(S3 Tables 内置)短暂可见
IceStream (第三方)独立服务接近实时需部署运维取决于配置
Flink RewriteDataFiles taskFlink 作业内Compaction 触发频率Compaction 前可见

9. 限制与后续计划

当前限制

后续计划

  1. 集成为 IcebergSink 的 inline post-commit task — 用户无需单独配置 maintenance job
  2. 自动清理 staging 分支已处理的 snapshots
  3. 添加官方文档
  4. DVResolver 支持 partition pruning 优化
  5. 保留 Row Lineage

10. 参考链接

#链接说明
1PR #15996ConvertEqualityDeletes 主 PR
2设计文档Google Doc 详细设计
3邮件列表讨论Iceberg Dev 邮件列表
4Data Lakehouse Weekly 2026-07-18社区报道
5The Equality Delete Problem问题深度分析 (RisingWave)
6Flink TableMaintenance 文档Iceberg 官方文档
7Issue #11122V3 Position Delete 改进 Umbrella Issue
8IceStream第三方 eq-to-DV 转换方案
9Streaming Updates in IcebergPosition Deletes at Scale (Etleap)
10Multi-Engine SupportIceberg 多引擎版本兼容矩阵

调研时间:2026-07-21 | 作者:Xiao Huang | 基于 PR #15996 及社区公开讨论整理