编程 Apache Iceberg V3 深度拆解:当删除向量取代删除文件——从 Puffin 字节级解析到行级血缘 CDC 的完整实战指南(2026)

2026-08-13 22:29:43

如果你在 2020 年问「数据湖用什么表格式」,答案是一场混战;如果你在 2026 年问同样的问题,答案基本已经收敛。但真正值得关心的不是「谁赢了」,而是 V3 规范里那几个看着不起眼的改动——它们把湖仓从「能用」推到了「能当主库用」。

这篇文章不讲「Iceberg 是什么」的科普。我假设你已经知道它是一个开放表格式,知道它有时间旅行和 Schema 演进。我要拆的是更硬的东西:V2 的删除机制到底烂在哪里、V3 的删除向量凭什么快一个数量级、行级血缘怎么把 CDC 从「重扫全表」变成「读一列」、以及这些东西落到生产环境会在哪些地方咬你一口。

全文配可运行代码,包括一段手写的 Puffin 删除向量二进制解析器——因为我发现绝大多数讲 Iceberg 的文章都停在「DV 用了 Roaring Bitmap」这一句,没人告诉你那个文件里的字节到底长什么样。


一、背景:表格式为什么会成为数据栈的地基

1.1 Hive 表格式的四宗罪

先说清楚我们在解决什么问题。Hive 表格式(我说的是「格式」,不是 Hive 这个引擎)的定义方式非常朴素:一个表就是一个目录,一个分区就是一个子目录,目录里的文件就是数据。

/warehouse/db/orders/
  dt=2026-08-13/
    part-00000.parquet
    part-00001.parquet
  dt=2026-08-12/
    ...

这套设计在 2009 年是天才,在 2026 年是灾难。四个根本缺陷:

第一,表的状态没有单一真相源。 「这个表现在有哪些文件」这个问题,只能靠 LIST 目录来回答。在 HDFS 上这是一次 NameNode RPC,在 S3 上这是一次最终一致性的分页列举——文件写完了但列举不到,是 S3 上无数数据事故的起点。

第二,没有原子性。 一个 INSERT OVERWRITE 要覆盖 100 个分区,中途失败了,你得到一个「一半新一半旧」的表。没有回滚,只能人工救。

第三,分区是物理耦合的。 分区字段必须以目录名的形式存在,查询必须显式带上分区谓词才能剪枝。用户写 WHERE order_time > '2026-08-01' 而分区键是 dt,全表扫描;改分区策略等于重写整张表。

第四,Schema 演进靠位置匹配。 Parquet 文件里的列靠名字或顺序对齐,删一列再加一列同名列,历史数据会读出脏值。

1.2 Iceberg 的破局思路:把文件列表变成元数据

Iceberg 的核心思想一句话能说完:不要靠列举目录来定义表,用一份显式的、不可变的、可版本化的元数据来定义表。

这一句话带来的连锁反应是巨大的:

  • 文件列表在元数据里 → 不需要 LIST,S3 一致性问题消失
  • 元数据不可变 + 指针原子切换 → 天然的快照隔离和时间旅行
  • 元数据里带每个文件的列统计 → 不靠目录名也能剪枝,分区可以「隐藏」
  • 元数据里带 field-id → Schema 演进靠 ID 不靠名字,改名不影响读

1.3 V1 → V2 → V3:三代规范的分水岭

规范版本核心能力解决的核心矛盾遗留问题
V1快照、隐藏分区、Schema 演进「表状态不可靠」只能整文件重写,无法行级更新
V2删除文件(position/equality delete)、Merge-on-Read「行级更新代价太高」删除文件读放大严重、无行标识
V3删除向量、行级血缘、Variant、地理类型、纳秒时间、默认值、多参数变换「读放大」+「无法做增量」+「半结构化」元数据提交仍是文件级

到 Iceberg Java 1.11.0(也就是当前 latest)为止,V3 已经从「规范定稿」走到了「主流引擎可用」。Snowflake 在 2026 年 4 月宣布对 V3 的支持,覆盖 Variant 半结构化类型、地理空间类型、行级血缘、删除向量、纳秒时间戳;DuckDB 的 iceberg 扩展在 1.5.3 里补齐了 MERGE INTO、ALTER TABLE 和 V3 支持。

这就是本文的重点:V3 不是一堆小特性的集合,它是让湖仓具备「主库能力」的最后一块拼图。


二、核心概念:三层元数据的完整解剖

要理解 V3 的改动,必须先精确理解元数据的层次。我见过太多人把这一层讲成「元数据 → 数据」两层,那样根本解释不了删除向量为什么能在查询规划阶段就生效。

2.1 四层指针链

        ┌──────────────┐
        │   Catalog    │  唯一可变的东西:表名 → 当前元数据文件路径
        └──────┬───────┘
               │  原子 CAS 交换(这是整个系统唯一的写锁点)
               ▼
        ┌──────────────────────────┐
        │ v42.metadata.json        │  表级:schemas / partition-specs /
        │                          │  sort-orders / snapshots / refs /
        │                          │  properties / next-row-id
        └──────┬───────────────────┘
               │  current-snapshot-id
               ▼
        ┌──────────────────────────┐
        │ snap-xxx.avro            │  Manifest List:本快照包含哪些 manifest
        │ (manifest list)          │  每条带分区上下界摘要 → 可跳过整个 manifest
        └──────┬───────────────────┘
               │
               ▼
        ┌──────────────────────────┐
        │ xxx-m0.avro (manifest)   │  文件级:每条一个 data file 或 delete file
        │                          │  带 record_count / 列 min-max / null 计数 /
        │                          │  split_offsets / first_row_id
        └──────┬───────────────────┘
               │
               ▼
        ┌──────────────────────────┐
        │ 00000-0-xxx.parquet      │  真实数据
        │ xxx.puffin               │  V3 删除向量
        └──────────────────────────┘

关键点:查询规划是自上而下逐层剪枝的。一个带分区谓词的查询,可能在 manifest list 层就砍掉 99% 的 manifest,根本不会去读那些 manifest 文件本身。这就是为什么 Iceberg 能在百万级文件的表上做到秒级规划——它的规划复杂度不是 O(文件数),而是 O(命中的 manifest 数)。

2.2 快照与原子提交:乐观并发的实现

Iceberg 的写入是乐观并发控制(OCC)。完整流程:

1. 读取当前 metadata(记住 base-snapshot-id 和 metadata 版本号)
2. 写数据文件到存储(此时对读者完全不可见)
3. 生成新的 manifest / manifest list
4. 生成新的 metadata.json(版本号 +1)
5. 向 Catalog 发起 CAS:
     "把 current-metadata 从 v42 换成 v43,前提是它现在还是 v42"
6a. 成功 → 提交完成,新快照对所有读者原子可见
6b. 失败(有人抢先提交)→ 校验冲突,可重试则回到步骤 3 重建元数据

第 6b 步的「校验冲突」是重点。Iceberg 不是无脑重试,它会根据隔离级别做断言:

  • snapshot 隔离:只要我改的文件没被别人删掉,就可以重放
  • serializable 隔离:如果别人在我的谓词范围内新增了文件,必须失败

配置项:

ALTER TABLE db.orders SET TBLPROPERTIES (
  'commit.retry.num-retries' = '10',       -- 默认 4,高并发场景建议调大
  'commit.retry.min-wait-ms' = '100',
  'commit.retry.max-wait-ms' = '60000',
  'commit.retry.total-timeout-ms' = '1800000',
  'write.merge.isolation-level' = 'snapshot'  -- 默认 serializable
);

这里有个非常常见的生产误区:很多人以为 Iceberg 支持高并发写入。准确说法是——它支持高并发的「不冲突」写入。所有写入最终都要争抢 Catalog 那一个 CAS 点。如果你有 50 个 Flink 任务每 10 秒 checkpoint 一次往同一张表写,你会看到大量提交冲突和重试风暴。解法是按分区拆表或者拆分支(branch),而不是调大重试次数。

2.3 序列号:删除文件如何知道自己该管谁

这是理解 MoR 的关键机制,也是最容易被讲错的地方。

每个快照有一个单调递增的 sequence number。每个 manifest 和其中的文件条目继承所在快照的序列号(这叫「继承式序列号」,好处是重用旧 manifest 时不用重写)。

删除文件的作用范围规则:

  • Position Delete:只作用于它显式引用的那个数据文件路径,且数据文件的序列号 <= 删除文件的序列号
  • Equality Delete:作用于所有序列号 严格小于 它的数据文件

第二条是 equality delete 的性能地狱之源:一个 equality delete 文件,理论上要和整张表里比它老的所有数据文件做 join。

时间轴 →
seq=1  data-A.parquet (1000 行)
seq=2  data-B.parquet (1000 行)
seq=3  eq-delete-1.parquet (删除 user_id IN (7, 42))
       ↑ 这个文件要作用于 data-A 和 data-B
seq=4  data-C.parquet (1000 行)
       ↑ 这个文件不受 eq-delete-1 影响(seq 更大)
seq=5  eq-delete-2.parquet
       ↑ 作用于 A, B, C

读一次全表:需要把 eq-delete-1 和 eq-delete-2 的键全部加载成哈希集合,然后对 A/B/C 每一行做探测。这是 O(数据行数 × 删除文件数) 的哈希探测。


三、V2 的技术债:删除文件为什么会失控

3.1 Position Delete 的三重代价

Position delete 文件的 schema 很简单:

file_path: string   -- 被删行所在的数据文件路径
pos: long           -- 被删行在该文件中的行号(从 0 开始)
row: struct         -- 可选,被删行的原始内容

看着很轻量,但生产上有三个致命问题:

代价一:一个数据文件可以被 N 个删除文件引用。

规范没有限制「一个数据文件只能有一个 position delete」。每次 MoR 写入都可能产生新的删除文件。所以真实场景是:

data-A.parquet
  ← pos-del-1.parquet (删了 pos 5, 17)
  ← pos-del-2.parquet (删了 pos 88)
  ← pos-del-3.parquet (删了 pos 5, 200)   -- 注意,可以重复!
  ← ... 累积到几十个

读 data-A 时,reader 必须打开所有引用它的删除文件,做归并排序,去重,才能得到「哪些行该跳过」的最终答案。每个删除文件都是一次独立的对象存储 GET(S3 上典型 20–100ms 首字节延迟),几十个删除文件意味着几秒钟的纯 I/O 等待,而这些等待完全无法被数据扫描的流水线掩盖。

代价二:内存放大。

Position delete 是以「行」为单位存储的。删掉 1000 万行,就是 1000 万条 (file_path, pos) 记录。file_path 是字符串,即使有字典编码,解码到内存里也是实打实的对象。一个删了 10% 数据的 TB 级表,光删除文件的内存开销就能把 executor 打爆。

代价三:查询规划阶段看不见删除量。

这一点最隐蔽也最要命。规划器知道 data-A 有 100 万行(record_count 在 manifest 里),但不知道其中有多少行已被删除。所以:

  • 统计信息 SELECT count(*) 无法走元数据快路径,必须真扫
  • CBO 的基数估算严重偏差 → join 顺序选错 → 查询慢 10 倍
  • split 切分按原始文件大小算,实际有效行数可能只剩 1%,导致大量空跑的 task

3.2 一个真实的数量级推演

假设一张 CDC 同步表:

  • 数据量 5TB,Parquet 文件 512MB 一个 → 约 10000 个数据文件
  • 每 5 分钟一次 upsert 微批,每批影响 2000 个数据文件
  • 一天 288 批 → 每个数据文件平均被 288 × 2000 / 10000 ≈ 58 个删除文件引用

一次全表扫描要打开 10000 + 58 × 10000 ≈ 59 万个文件。即使每个文件只花 30ms 建立连接,串行化的部分也足以让查询从「分钟级」退化到「小时级」。这就是业界俗称的 delete file explosion(删除文件爆炸)。

在 V2 下,唯一的解法是高频压实:跑 rewrite_position_delete_files 和 rewrite_data_files,把删除折进数据里。但压实本身要重写数据,成本极高,而且压实任务和写入任务会互相冲突。这是一个用运维成本硬扛架构缺陷的典型例子。


四、V3 核心特性逐个拆解

4.1 删除向量:从「行列表」到「位图」

V3 的删除向量(Deletion Vector, DV)做了两个决定性改变:

改变一(规范级约束):一个数据文件在一个快照里最多只能有一个 DV。

这一条比任何性能优化都重要。它把「读 N 个删除文件并归并」变成了「读 1 个 DV」。写入方有义务在提交时把旧 DV 和新删除合并成一个新 DV——把复杂度从读路径挪到了写路径。读多写少的分析场景,这笔账怎么算都赚。

改变二(编码级优化):用 Roaring Bitmap 存行位置,而不是行记录。

Roaring Bitmap 把 32 位整数空间切成 65536 个「桶」,每个桶根据基数自适应选择三种容器:

容器类型触发条件存储方式单值成本
Array Container基数 ≤ 4096排序的 uint16 数组2 字节
Bitmap Container基数 > 40968KB 定长位图~0.125 字节
Run Container存在长连续段(起点, 长度) 对列表摊薄到接近 0

对删除场景这简直是量身定制:

  • 稀疏删除(随机删几百行)→ Array Container,几百字节
  • 稠密删除(删一大段,比如按时间范围删)→ Run Container,几十字节
  • 中等删除 → Bitmap Container,8KB 封顶

Iceberg 用的是 64 位 Roaring(因为行位置是 long),实现上是「高 32 位分桶 → 每桶一个 32 位 Roaring Bitmap」的 portable 格式。

内存对比(删除 100 万行):

方案存储大小内存中形态单行判定复杂度
Position Delete~10–30MB(Parquet 压缩后)100 万条记录对象归并排序 O(log N × 文件数)
Deletion Vector~125KB(Bitmap Container)常驻位图O(1) 位测试

两个数量级的差距,而且是「存储 + 内存 + CPU」三重收益。

4.2 Puffin 文件格式:DV 的物理载体

DV 不是裸的位图文件,它装在 Puffin 容器里。Puffin 是 Iceberg 自己定义的「辅助统计信息文件格式」,同一个 Puffin 文件可以装多个 blob(多个数据文件的 DV,或者 Theta Sketch 之类的基数草图)。

Puffin 文件的物理布局:

┌────────────────────────────────────────┐
│ Magic: "PFA1" (4 字节: 0x50 0x46 0x41 0x31) │
├────────────────────────────────────────┤
│ Blob 1 payload                          │
│ Blob 2 payload                          │
│ ...                                     │
├────────────────────────────────────────┤
│ Footer:                                 │
│   Magic "PFA1"                          │
│   FooterPayload (JSON, 可选压缩)         │
│     { "blobs": [                        │
│         { "type": "deletion-vector-v1", │
│           "fields": [],                 │
│           "snapshot-id": ...,           │
│           "sequence-number": ...,       │
│           "offset": 4,                  │
│           "length": 1234,               │
│           "properties": {               │
│              "cardinality": "17"        │
│           } } ] }                       │
│   FooterPayloadSize (4 字节 LE)          │
│   Flags (4 字节)                         │
│   Magic "PFA1"                          │
└────────────────────────────────────────┘

而 deletion-vector-v1 这个 blob 的内部布局是:

┌──────────────────────────────────────────────┐
│ length: 4 字节 big-endian                     │  = 4(magic) + bitmap 字节数
├──────────────────────────────────────────────┤
│ magic: 0xD1D33964 (小端写入,字节序 64 39 D3 D1) │
├──────────────────────────────────────────────┤
│ Roaring Bitmap (portable 64-bit 格式)         │
├──────────────────────────────────────────────┤
│ CRC-32: 4 字节 big-endian                     │  校验 magic + bitmap
└──────────────────────────────────────────────┘

**为什么要在 blob 内部再套一层 length + magic + CRC?**因为 Puffin 文件是共享的——多个数据文件的 DV 可以放在同一个 Puffin 里。Manifest 里记录的是 (puffin 路径, content_offset, content_size_in_bytes) 三元组,reader 可以直接发一个 Range GET 精准取出那一段字节,不需要解析整个 Puffin footer。而 magic + CRC 保证了「我 range 读到的这一段确实是我要的那个 DV」——这是对象存储上防止读到脏数据的必要保险。

Manifest 里 DV 条目新增的字段:

content = 1 (POSITION_DELETES)
file_path = "s3://bucket/data/xxx.puffin"
referenced_data_file = "s3://bucket/data/00000-0-abc.parquet"   ← V3 新增,必填
content_offset = 4                                              ← V3 新增
content_size_in_bytes = 1234                                    ← V3 新增
record_count = 17                                               ← 被删行数,规划期可用!

record_count 这一个字段就解决了 3.1 节的「代价三」:规划器现在能算出每个数据文件的有效行数 = record_count(data) − record_count(dv)。 COUNT(*) 重回元数据快路径,CBO 基数估算恢复准确,split 切分不再空跑。

4.3 行级血缘:让 CDC 从「重扫」变成「读列」

这是 V3 里我个人认为最被低估的特性。

问题背景:湖仓上做增量处理有个死结。你想知道「上一个快照到这个快照之间,哪些行变了」。Iceberg V2 提供了 changelog 视图,但它的实现是:对比两个快照的文件集合,把新增文件全读一遍,把删除文件全读一遍,然后推断出 insert/delete/update_before/update_after。

这个方案的致命缺陷是它无法区分「真实的数据变更」和「压实产生的物理重写」。你跑一次 rewrite_data_files 把 1000 个小文件压成 10 个大文件,数据一行没变,但 changelog 会告诉你「1000 个文件的所有行都被删了,10 个文件的所有行都是新增的」。下游全量重算。

V3 的解法:给每一行一个稳定的、跨物理重写不变的身份。

两个保留字段:

  • _row_id:行的稳定唯一标识
  • _last_updated_sequence_number:这一行最后一次被修改时的序列号

关键在于它的赋值机制——不是给每行都物理存一个 ID,那样太浪费。 机制是这样的:

表元数据里有一个全局计数器: next-row-id

提交新快照时:
  1. 为本次新增的每个 manifest 分配 first_row_id
  2. 每个 data file 条目记录自己的 first_row_id
  3. next-row-id += 本次新增的总行数

读取时:
  某行的 _row_id = 如果文件里物理存了 → 用物理值
                   否则 → data_file.first_row_id + 行在文件内的位置

这就是「隐式分配 + 显式覆盖」的混合策略。新写入的行零存储开销(只在文件元数据里多一个 long);压实重写的行则物理保留原有的 _row_id,从而身份不变。

同理,_last_updated_sequence_number 的规则是:

  • 新写入的行 → 继承数据文件的 data_sequence_number
  • 压实搬运的行 → 保留原值(这一点是压实器的强制义务)

落到实践上的价值:

-- V2 时代做增量:全量对比文件集,压实会污染结果
-- V3 时代做增量:一个谓词搞定
SELECT *, _row_id, _last_updated_sequence_number
FROM db.orders
WHERE _last_updated_sequence_number > 12345;  -- 上次处理到的水位

而且因为 _last_updated_sequence_number 是数据文件级别的元数据(同一文件内多数为常量),这个谓词可以在 manifest 层就完成文件剪枝——压实产生的文件天然带着旧序列号,直接被跳过。

这才是「湖仓能做流处理」的真正基础。物化视图的增量刷新、下游数仓的 CDC 同步、特征平台的增量回填,全都建立在这个能力上。

4.4 Variant:半结构化数据的正确解法

在 V3 之前,湖上存 JSON 只有两条路,都很难受:

路子 A:存成 string,查询时 json_extract。 每次查询都要 parse 整个 JSON 文本,无法剪枝,无法用统计信息,CPU 全烧在解析上。

路子 B:展开成 struct。 需要预先知道所有字段,schema 一变就要改表;字段几千个的埋点数据直接爆炸。

Variant 是第三条路:二进制自描述编码 + 可选的列裁剪(shredding)。

Variant 值由两部分组成:

  • metadata:一个字典,存这个 variant 里出现的所有 key 字符串(去重)
  • value:递归的二进制编码,对象里的 key 用「字典下标」引用,不重复存字符串

好处:

  1. 不需要 parse 文本。取 payload.user.id 是二进制偏移跳转,不是字符串扫描
  2. key 字符串只存一次。埋点数据里 "event_timestamp" 这种长 key 出现 100 万次,只在 metadata 里存一份
  3. 保留类型。JSON 里的 123 和 "123" 是不同的,round-trip 不丢信息
  4. Schema 无关。写入方随便加字段,表结构不用动

Shredding(列裁剪)是性能关键。 如果你知道 payload.user_id 和 payload.event_type 是高频查询字段,可以让写入方把这两个字段同时以独立的 typed column 形式物化出来。这样:

  • 查 payload.user_id = 42 → 走物化列,有 min-max 统计,能剪枝,能向量化
  • 查 payload.some_rare_field → 回落到 variant 二进制解析

这本质上是「行存的灵活性 + 列存的性能」的组合拳,代价是写入端多一点 CPU 和存储。

4.5 其他 V3 特性快速过一遍

地理空间类型(geometry / geography)

  • geometry:平面坐标,WKB 编码,带 CRS 参数(默认 OGC:CRS84)
  • geography:球面坐标,额外带边插值算法(球面大圆等)
  • 统计信息支持地理边界框 → 空间查询可以做分区/文件剪枝

意义:以前做地理分析必须上 PostGIS 或专门的空间数据库,现在湖上原生支持,且能跨引擎互操作。

纳秒时间戳(timestamp_ns / timestamptz_ns)

微秒精度在金融高频、可观测性(span 时间)、IoT 场景不够用。V3 补上纳秒。注意坑:V3 的 timestamp_ns 在有效范围上和微秒版不一样(纳秒 long 只能表示约 ±292 年),跨版本转换要小心溢出。

unknown 类型

一个「永远为 null」的类型。听着没用,实际是 Schema 演进的重要拼图:你可以先加一个 unknown 列占位(比如上游还没定类型),后续再演进成具体类型,而不用经历「删列再加列」这种会污染历史数据的操作。

列默认值(initial-default / write-default)

ALTER TABLE db.orders ADD COLUMN region string DEFAULT 'unknown';

V3 之前,加列的默认值只能是 null(因为历史 Parquet 文件里根本没这列)。V3 把默认值拆成两个:

  • initial-default:读到没有这列的老文件时,填这个值
  • write-default:新写入不指定该列时,填这个值

这意味着「加带默认值的列」变成 O(1) 元数据操作,不需要重写 PB 级历史数据。 从关系数据库迁过来的人会对这个特性非常有感——这是 ALTER TABLE ADD COLUMN ... DEFAULT 在湖上的第一次真正可用。

多参数分区/排序变换

V3 允许变换函数接受多个源列。典型用途是复合键的分桶:以前只能 bucket(16, user_id),现在可以对多列组合做变换,让相关数据物理聚集在一起。对 join 场景(storage-partitioned join)价值很大。


五、架构分析:Copy-on-Write 与 Merge-on-Read 的决策模型

V3 让 MoR 变得便宜了,但不代表 MoR 永远正确。这里给一个可以直接套用的量化模型。

5.1 两种模式的成本结构

Copy-on-Write(CoW):更新一行 = 把这行所在的整个数据文件读出来、改掉、重写成新文件。

写成本 = Σ(被触碰文件的大小) / 变更行的数据量  = 写放大倍数
读成本 = 0 额外开销(数据永远是「干净」的)

假设文件 512MB,平均每行 1KB(即约 52 万行/文件),你更新了 1 行:

写放大 = 512MB / 1KB = 524288 倍

五十万倍。 这就是为什么 CoW 绝对不能用于高频点更新。

Merge-on-Read(MoR):更新一行 = 写一条新数据 + 在 DV 里标记老行删除。

写成本 ≈ O(1)(一条新记录 + 一个位翻转)
读成本 = 每个数据文件额外一次 DV Range GET + 位图判定

5.2 决策公式

定义:

  • u = 单位时间的更新行数
  • r = 单位时间的全表扫描次数
  • F = 数据文件数
  • S = 平均文件大小
  • L = 对象存储单次请求延迟(S3 典型 30–80ms)
  • T = 压实周期

CoW 的单位时间成本(以字节 I/O 计):

Cost_cow ≈ u × S × (被触碰文件的离散度系数 k)

其中 k ∈ [1/rows_per_file, 1],更新越随机 k 越接近 1(每次更新都碰一个不同的文件)。

MoR 的单位时间成本:

Cost_mor ≈ u × row_size          (写)
         + r × F × L / 并发度      (读时的 DV 请求开销)
         + 压实成本 / T

交叉点:当 u × S × k > r × F × L,选 MoR;反之选 CoW。

翻译成人话的经验法则:

场景特征推荐模式理由
日更新率 < 1%,读频极高CoW写放大可承受,换零读开销
CDC 同步、每分钟级 upsertMoR写放大不可承受
GDPR 删除(按用户删,极稀疏)MoR + DVDV 对稀疏删除是最优编码
按时间范围批量删(连续段)MoR + DVRun Container,几乎零成本
更新集中在最近分区分区级混合热分区 MoR,冷分区 CoW

最后一行是生产上真正的答案。 Iceberg 支持按分区设置不同策略——虽然表属性是表级的,但你可以让压实任务只对冷分区做「转 CoW」的重写,热分区保持 MoR。

5.3 配置层面怎么落

ALTER TABLE db.orders SET TBLPROPERTIES (
  -- 三种 DML 分别配置,注意 delete/update/merge 是独立的
  'write.delete.mode' = 'merge-on-read',
  'write.update.mode' = 'merge-on-read',
  'write.merge.mode'  = 'merge-on-read',

  -- V3 关键:删除粒度用 file(DV),而不是 partition
  'write.delete.granularity' = 'file',

  -- 目标文件大小,默认 512MB
  'write.target-file-size-bytes' = '536870912',

  -- 写入分布模式:hash 能避免小文件,代价是一次 shuffle
  'write.distribution-mode' = 'hash',

  -- 元数据文件清理,不设会无限累积 metadata.json
  'write.metadata.delete-after-commit.enabled' = 'true',
  'write.metadata.previous-versions-max' = '20'
);

write.delete.granularity 这个参数值得单独说。它有两个值:

  • partition:一个删除文件覆盖整个分区内的多个数据文件 → 删除文件少但读时要过滤
  • file:一个删除文件(DV)只对应一个数据文件 → 这是 V3 DV 的正确配置

V3 下必须用 file,否则你拿不到 DV 的核心收益(一对一映射 + range GET)。


六、代码实战

下面的代码都是可运行的。环境假设:Spark 4.x + Iceberg 1.11.0,Python 3.11。

6.1 环境准备与建 V3 表

# 启动带 Iceberg 的 Spark SQL,注意 format-version=3
spark-sql \
  --packages org.apache.iceberg:iceberg-spark-runtime-4.0_2.13:1.11.0 \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
  --conf spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog \
  --conf spark.sql.catalog.lake.type=rest \
  --conf spark.sql.catalog.lake.uri=http://localhost:8181 \
  --conf spark.sql.catalog.lake.warehouse=s3://my-lake/wh \
  --conf spark.sql.catalog.lake.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
  --conf spark.sql.defaultCatalog=lake
CREATE TABLE lake.ops.events (
  event_id      bigint,
  user_id       bigint,
  event_type    string,
  -- V3 纳秒时间戳
  occurred_at   timestamp_ns,
  -- V3 Variant:埋点自由字段
  payload       variant,
  -- V3 地理类型
  location      geography,
  -- V3 列默认值:加列不重写历史数据
  channel       string DEFAULT 'unknown'
)
USING iceberg
PARTITIONED BY (days(occurred_at), bucket(16, user_id))
TBLPROPERTIES (
  'format-version' = '3',
  'write.delete.mode' = 'merge-on-read',
  'write.update.mode' = 'merge-on-read',
  'write.merge.mode'  = 'merge-on-read',
  'write.delete.granularity' = 'file',
  -- V3 行级血缘(V3 表默认开启,显式写出来更清楚)
  'write.row-lineage' = 'true',
  'write.parquet.compression-codec' = 'zstd',
  'write.parquet.compression-level' = '3'
);

几个决策解释:

  • days(occurred_at) 而不是 hours(...):分区粒度太细会导致小文件和 manifest 膨胀。经验值是让每个分区至少有 1GB 数据
  • bucket(16, user_id):二级分桶让按 user_id 的点查能剪掉 15/16 的文件,同时为 storage-partitioned join 铺路
  • zstd level 3:在压缩率和 CPU 之间的最佳平衡点。别用 level 9,收益 <5% 但 CPU 翻倍

6.2 MERGE INTO:upsert 的标准写法

-- 假设 stg.event_changes 是 CDC 落下来的变更集
MERGE INTO lake.ops.events t
USING (
  -- 关键:源表必须先去重,否则「同一个 key 匹配多行」会直接报错
  SELECT * FROM (
    SELECT *, ROW_NUMBER() OVER (
      PARTITION BY event_id ORDER BY cdc_seq DESC
    ) AS rn
    FROM stg.event_changes
  ) WHERE rn = 1
) s
ON t.event_id = s.event_id
   -- 强烈建议带上分区谓词,让 Iceberg 只扫相关分区
   AND t.occurred_at >= timestamp_ns '2026-08-01 00:00:00'
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED AND s.op = 'U' THEN UPDATE SET
     t.event_type = s.event_type,
     t.payload    = s.payload,
     t.channel    = s.channel
WHEN NOT MATCHED AND s.op != 'D' THEN INSERT *;

ON 条件里那行分区谓词是性能生死线。 没有它,Iceberg 必须假设匹配可能发生在任何分区,于是扫全表。加上它,规划期就砍掉 99% 的文件。我见过的 MERGE 慢案例里,八成是这个原因。

6.3 用 PyIceberg 直接读元数据:看清 DV 到底在哪

"""
用 pyiceberg 检查表的物理布局,重点是 V3 删除向量的元数据。
pip install "pyiceberg[s3fs,pyarrow]>=0.9"
"""
from pyiceberg.catalog import load_catalog

catalog = load_catalog(
    "lake",
    **{
        "type": "rest",
        "uri": "http://localhost:8181",
        "warehouse": "s3://my-lake/wh",
        "s3.endpoint": "http://localhost:9000",
    },
)

tbl = catalog.load_table("ops.events")

print("=" * 60)
print(f"format-version : {tbl.metadata.format_version}")
print(f"current snapshot: {tbl.metadata.current_snapshot_id}")
print(f"next-row-id    : {getattr(tbl.metadata, 'next_row_id', 'N/A')}")
print("=" * 60)

# 1) 快照历史:看清每次提交干了什么
print("\n[快照历史]")
for snap in tbl.metadata.snapshots:
    summary = snap.summary or {}
    print(
        f"  seq={snap.sequence_number:<5} id={snap.snapshot_id} "
        f"op={summary.get('operation'):<9} "
        f"+files={summary.get('added-data-files', 0):<5} "
        f"+dv={summary.get('added-position-delete-files', 0):<4} "
        f"+rows={summary.get('added-records', 0)}"
    )

# 2) 数据文件清单 + 有效行数(V3 才能在规划期算出来)
print("\n[数据文件与有效行数]")
files_df = tbl.inspect.files().to_pylist()
total_raw, total_deleted = 0, 0
for f in files_df[:10]:
    raw = f["record_count"]
    total_raw += raw
    print(
        f"  {f['file_path'].split('/')[-1][:44]:<46} "
        f"rows={raw:<9} size={f['file_size_in_bytes'] / 1024 / 1024:.1f}MB "
        f"first_row_id={f.get('first_row_id')}"
    )

# 3) 删除文件(V3 里就是 Puffin DV)
print("\n[删除向量]")
dv_df = tbl.inspect.delete_files().to_pylist()
for d in dv_df[:10]:
    total_deleted += d["record_count"]
    print(
        f"  puffin={d['file_path'].split('/')[-1][:36]:<38} "
        f"refs={str(d.get('referenced_data_file', '')).split('/')[-1][:30]:<32} "
        f"offset={d.get('content_offset')} "
        f"size={d.get('content_size_in_bytes')} "
        f"deleted_rows={d['record_count']}"
    )

print(f"\n原始行数={total_raw}  已删行数={total_deleted}  "
      f"有效行数={total_raw - total_deleted}")
print(f"删除率={total_deleted / max(total_raw, 1) * 100:.2f}%  "
      f"(>15% 该压实了)")

这段脚本我在生产上是当巡检工具用的。最后那个「删除率」是决定是否触发压实的核心指标。

6.4 手写 Puffin 删除向量解析器

这是本文我最想给你的东西。市面上所有讲 DV 的文章都停在概念层,我们直接把字节撕开看。

"""
puffin_dv_reader.py
纯 Python 解析 Iceberg V3 的 Puffin 删除向量文件。
不依赖任何 Iceberg 库,只靠 struct + 一个 Roaring Bitmap 解析器。

用途:
  1. 排障——怀疑 DV 有问题时,直接验证字节
  2. 教学——理解 DV 的真实物理形态
  3. 审计——不启动 Spark 就能算出某文件删了哪些行
"""
import json
import struct
import zlib
from dataclasses import dataclass
from typing import List, Set

PUFFIN_MAGIC = b"PFA1"
DV_BLOB_MAGIC = 0xD1D33964  # deletion-vector-v1 的 magic
FOOTER_STRUCT_LEN = 4 + 4 + 4  # payload_size(4) + flags(4) + magic(4)


@dataclass
class BlobMetadata:
    type: str
    offset: int
    length: int
    snapshot_id: int
    sequence_number: int
    properties: dict


def read_puffin_footer(data: bytes) -> List[BlobMetadata]:
    """解析 Puffin 尾部,拿到所有 blob 的元数据。"""
    if data[:4] != PUFFIN_MAGIC:
        raise ValueError(f"不是 Puffin 文件,头部魔数是 {data[:4]!r}")
    if data[-4:] != PUFFIN_MAGIC:
        raise ValueError("Puffin 尾部魔数缺失,文件可能被截断")

    # 从尾部往前定位 footer payload
    tail = data[-FOOTER_STRUCT_LEN:]
    payload_size, flags = struct.unpack("<II", tail[:8])

    # flags 的 bit 0 表示 footer payload 是否用 LZ4 压缩
    footer_compressed = bool(flags & 0x1)

    payload_end = len(data) - FOOTER_STRUCT_LEN
    payload_start = payload_end - payload_size
    payload = data[payload_start:payload_end]

    if footer_compressed:
        raise NotImplementedError(
            "footer 用了 LZ4 压缩,需要 pip install lz4 后用 lz4.block.decompress"
        )

    meta = json.loads(payload.decode("utf-8"))
    blobs = []
    for b in meta.get("blobs", []):
        blobs.append(
            BlobMetadata(
                type=b["type"],
                offset=b["offset"],
                length=b["length"],
                snapshot_id=b.get("snapshot-id", -1),
                sequence_number=b.get("sequence-number", -1),
                properties=b.get("properties", {}),
            )
        )
    return blobs


def parse_roaring_32(buf: bytes, pos: int) -> tuple[Set[int], int]:
    """
    解析一个标准 32 位 Roaring Bitmap(portable 格式)。
    返回 (值集合, 新的读取位置)。

    portable 格式布局:
      cookie(4)
      如果 cookie == 0x3B4C0000 | (n-1): 有 run 容器,跟 run bitmap
      如果 cookie == 0x00003B4C: 无 run,跟 size(4)
      然后是 keyscard 数组: (key uint16, cardinality-1 uint16) × n
      如果 n >= 4 或有 run: offset 数组 (uint32 × n)
      然后是各容器数据
    """
    SERIAL_COOKIE_NO_RUNCONTAINER = 0x00003B4C
    SERIAL_COOKIE = 0x3B4C

    start = pos
    cookie = struct.unpack_from("<I", buf, pos)[0]
    pos += 4

    has_run = False
    if cookie == SERIAL_COOKIE_NO_RUNCONTAINER:
        n_containers = struct.unpack_from("<I", buf, pos)[0]
        pos += 4
    elif (cookie & 0xFFFF) == SERIAL_COOKIE:
        n_containers = (cookie >> 16) + 1
        has_run = True
    else:
        raise ValueError(f"未知 Roaring cookie: 0x{cookie:08X}")

    run_flags = bytearray()
    if has_run:
        nbytes = (n_containers + 7) // 8
        run_flags = bytearray(buf[pos:pos + nbytes])
        pos += nbytes

    keys, cards = [], []
    for _ in range(n_containers):
        k, c = struct.unpack_from("<HH", buf, pos)
        pos += 4
        keys.append(k)
        cards.append(c + 1)  # 存的是 cardinality - 1

    # offset 数组:无 run 且容器数 >= 4 时存在(NO_RUN 情况下总是存在)
    if not has_run or n_containers >= 4:
        pos += 4 * n_containers  # 我们顺序读,不需要 offset

    result: Set[int] = set()
    for i in range(n_containers):
        high = keys[i] << 16
        card = cards[i]
        is_run = has_run and (run_flags[i // 8] & (1 << (i % 8))) != 0

        if is_run:
            # Run Container: n_runs(uint16) + (start uint16, len uint16) × n
            n_runs = struct.unpack_from("<H", buf, pos)[0]
            pos += 2
            for _ in range(n_runs):
                run_start, run_len = struct.unpack_from("<HH", buf, pos)
                pos += 4
                # 注意:存的 length 是「额外长度」,实际长度 = len + 1
                for v in range(run_start, run_start + run_len + 1):
                    result.add(high | v)
        elif card <= 4096:
            # Array Container: uint16 × card
            vals = struct.unpack_from(f"<{card}H", buf, pos)
            pos += 2 * card
            for v in vals:
                result.add(high | v)
        else:
            # Bitmap Container: 固定 8192 字节 = 1024 个 uint64
            words = struct.unpack_from("<1024Q", buf, pos)
            pos += 8192
            for wi, w in enumerate(words):
                if w == 0:
                    continue
                for bi in range(64):
                    if w & (1 << bi):
                        result.add(high | (wi * 64 + bi))

    return result, pos


def parse_deletion_vector(blob_payload: bytes) -> Set[int]:
    """
    解析 deletion-vector-v1 blob,返回被删除的行位置集合。

    布局: length(4, BE) | magic(4) | roaring_64 | crc32(4, BE)
    """
    if len(blob_payload) < 12:
        raise ValueError("DV blob 太短")

    declared_len = struct.unpack_from(">I", blob_payload, 0)[0]
    body_start = 4
    body_end = body_start + declared_len          # magic + bitmap
    crc_expected = struct.unpack_from(">I", blob_payload, body_end)[0]

    body = blob_payload[body_start:body_end]
    crc_actual = zlib.crc32(body) & 0xFFFFFFFF
    if crc_actual != crc_expected:
        raise ValueError(
            f"CRC 校验失败: 期望 0x{crc_expected:08X}, 实际 0x{crc_actual:08X}。"
            f"文件可能损坏,或者你 range GET 的偏移错了。"
        )

    magic = struct.unpack_from("<I", body, 0)[0]
    if magic != DV_BLOB_MAGIC:
        raise ValueError(f"DV magic 不匹配: 0x{magic:08X}, 期望 0x{DV_BLOB_MAGIC:08X}")

    # 64 位 Roaring: n_buckets(uint64 LE) + (high32 uint32, roaring32) × n
    pos = 4
    n_buckets = struct.unpack_from("<Q", body, pos)[0]
    pos += 8

    positions: Set[int] = set()
    for _ in range(n_buckets):
        high32 = struct.unpack_from("<I", body, pos)[0]
        pos += 4
        low_vals, pos = parse_roaring_32(body, pos)
        base = high32 << 32
        for v in low_vals:
            positions.add(base | v)

    return positions


def inspect_puffin(path: str) -> None:
    with open(path, "rb") as f:
        data = f.read()

    print(f"文件: {path}  大小: {len(data)} 字节")
    blobs = read_puffin_footer(data)
    print(f"包含 {len(blobs)} 个 blob\n")

    for i, b in enumerate(blobs):
        print(f"--- blob #{i} ---")
        print(f"  type       : {b.type}")
        print(f"  offset/len : {b.offset} / {b.length}")
        print(f"  snapshot   : {b.snapshot_id} (seq={b.sequence_number})")
        print(f"  properties : {b.properties}")

        if b.type != "deletion-vector-v1":
            print("  (非 DV,跳过解析)\n")
            continue

        payload = data[b.offset:b.offset + b.length]
        positions = parse_deletion_vector(payload)
        declared = int(b.properties.get("cardinality", -1))

        print(f"  解析出被删行数: {len(positions)}")
        if declared >= 0 and declared != len(positions):
            print(f"  ⚠️ 与 properties.cardinality({declared}) 不一致!元数据可能损坏")
        srt = sorted(positions)
        print(f"  前 20 个位置: {srt[:20]}")
        if srt:
            span = srt[-1] - srt[0] + 1
            density = len(srt) / span
            container = ("Run(连续段)" if density > 0.9
                         else "Bitmap(稠密)" if len(srt) > 4096
                         else "Array(稀疏)")
            print(f"  位置跨度: {srt[0]}~{srt[-1]}  密度={density:.4f}  "
                  f"推测容器类型={container}")
            print(f"  每行平均字节成本: {b.length / len(srt):.3f}")
        print()


if __name__ == "__main__":
    import sys
    inspect_puffin(sys.argv[1])

跑起来的输出大概是这样:

文件: 00000-dv-abc123.puffin  大小: 412 字节
包含 1 个 blob

--- blob #0 ---
  type       : deletion-vector-v1
  offset/len : 4 / 386
  snapshot   : 5847392847362 (seq=17)
  properties : {'cardinality': '143', 'referenced-data-file': 's3://...'}
  解析出被删行数: 143
  前 20 个位置: [12, 88, 401, 1203, 1204, 1205, ...]
  位置跨度: 12~524288  密度=0.0003  推测容器类型=Array(稀疏)
  每行平均字节成本: 2.699

注意最后一行:每删一行平均 2.7 字节。 对比 position delete 方案里每行至少几十字节(file_path 字符串 + long),差距一目了然。而如果删除是连续段(比如按时间范围删),Run Container 能把这个数字压到 0.1 字节以下。

6.5 用行级血缘做增量 CDC 管道

"""
incremental_cdc.py
基于 V3 行级血缘的增量消费管道。
核心优势:压实(compaction)不会产生假变更。
"""
import json
from pathlib import Path
from pyspark.sql import SparkSession
from pyspark.sql import functions as F

WATERMARK_FILE = Path("/var/lib/cdc/events_watermark.json")


def load_watermark() -> int:
    if WATERMARK_FILE.exists():
        return json.loads(WATERMARK_FILE.read_text())["last_seq"]
    return 0


def save_watermark(seq: int) -> None:
    WATERMARK_FILE.parent.mkdir(parents=True, exist_ok=True)
    WATERMARK_FILE.write_text(json.dumps({"last_seq": seq}))


def run_batch(spark: SparkSession) -> None:
    last_seq = load_watermark()

    # 1) 当前最新序列号(作为本批的上界,保证幂等和可重放)
    current_seq = (
        spark.sql("SELECT max(sequence_number) AS s FROM lake.ops.events.snapshots")
        .collect()[0]["s"]
    )
    if current_seq is None or current_seq <= last_seq:
        print(f"无新数据(水位 {last_seq}),跳过")
        return

    print(f"处理区间: ({last_seq}, {current_seq}]")

    # 2) 关键查询:只读被修改过的行
    #    _last_updated_sequence_number 大多在文件级是常量,
    #    因此这个谓词能在 manifest 层剪掉绝大部分文件
    changed = (
        spark.read.table("lake.ops.events")
        .select(
            "_row_id",
            "_last_updated_sequence_number",
            "event_id",
            "user_id",
            "event_type",
            "occurred_at",
        )
        .where(
            (F.col("_last_updated_sequence_number") > last_seq)
            & (F.col("_last_updated_sequence_number") <= current_seq)
        )
    )

    n = changed.count()
    print(f"变更行数: {n}")

    if n > 0:
        # 3) _row_id 是稳定标识 → 下游可以直接用它做幂等 upsert
        (
            changed.write.format("jdbc")
            .option("url", "jdbc:postgresql://dw:5432/analytics")
            .option("dbtable", "stg_events_delta")
            .option("truncate", "true")
            .mode("overwrite")
            .save()
        )
        # 下游 SQL:
        #   INSERT INTO dim_events SELECT * FROM stg_events_delta
        #   ON CONFLICT (row_id) DO UPDATE SET ...

    # 4) 全部成功后才推进水位(失败重跑不丢数据)
    save_watermark(current_seq)
    print(f"水位推进到 {current_seq}")


if __name__ == "__main__":
    spark = (
        SparkSession.builder.appName("incremental-cdc")
        .config(
            "spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
        )
        .getOrCreate()
    )
    run_batch(spark)

把这段代码和 V2 时代的方案对比一下,你就明白 row lineage 的价值:

维度V2 changelog 方案V3 row lineage 方案
压实是否产生假变更会,全表都变成「删+插」不会,序列号保留
需要读取的数据量快照差集的全部文件仅命中谓词的文件
下游幂等键需自己定义业务主键现成的 _row_id
能否文件级剪枝不能能(序列号在元数据里)
// Flink SQL 方式:往 V3 表做流式 upsert
// 提交: ./bin/sql-client.sh -j iceberg-flink-runtime-1.20-1.11.0.jar

// 1) 注册 catalog
CREATE CATALOG lake WITH (
  'type' = 'iceberg',
  'catalog-type' = 'rest',
  'uri' = 'http://localhost:8181',
  'warehouse' = 's3://my-lake/wh'
);

// 2) checkpoint 是提交粒度,这个值直接决定小文件数量
SET 'execution.checkpointing.interval' = '60s';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';

// 3) 源:Kafka Debezium CDC
CREATE TABLE kafka_events (
  event_id BIGINT,
  user_id  BIGINT,
  event_type STRING,
  occurred_at TIMESTAMP(9),
  PRIMARY KEY (event_id) NOT ENFORCED
) WITH (
  'connector' = 'kafka',
  'topic' = 'db.public.events',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'debezium-json',
  'scan.startup.mode' = 'group-offsets'
);

// 4) 写入:开启 upsert 模式
INSERT INTO lake.ops.events /*+ OPTIONS(
  'upsert-enabled' = 'true',
  'write.distribution-mode' = 'hash',
  'write.delete.granularity' = 'file',
  'compaction-enabled' = 'true',
  'compaction.small-file-threshold' = '134217728'
) */
SELECT event_id, user_id, event_type, occurred_at,
       CAST(NULL AS STRING) AS channel
FROM kafka_events;

Flink 写 Iceberg 的三个必踩坑,提前说:

  1. checkpoint 间隔就是提交间隔。 设 10s 你一天会产生 8640 个快照和天量小文件。低于 60s 我基本不建议,除非你有配套的实时压实。

  2. upsert-enabled=true 要求表有主键,且必须配 write.distribution-mode=hash。 否则同一个 key 的不同版本会落到不同的 writer subtask,导致删除标记和数据错位——这类数据不一致极难排查。

  3. Flink 的 equality delete 在 V3 下不会自动变成 DV。 DV 对应的是 position delete 语义。如果你的 Flink 版本还在写 equality delete,V3 的读放大收益你是拿不到的——升级 connector 并确认 write.delete.granularity=file 生效。

6.7 DuckDB 侧的轻量查询

不是所有查询都要开 Spark 集群。DuckDB 的 iceberg 扩展在 1.5.3 之后已经能读 V3 表,还支持 MERGE INTO 和 ALTER TABLE。

-- 一行装扩展
INSTALL iceberg;
LOAD iceberg;

-- 挂载 REST Catalog
ATTACH 'warehouse' AS lake (
  TYPE iceberg,
  ENDPOINT 'http://localhost:8181'
);

-- 直接查,本地单机就能跑
SELECT
  event_type,
  count(*) AS cnt,
  count(DISTINCT user_id) AS uv
FROM lake.ops.events
WHERE occurred_at >= '2026-08-01'
GROUP BY 1
ORDER BY cnt DESC;

-- 时间旅行
SELECT count(*) FROM lake.ops.events AT (VERSION => 5847392847362);

-- 看元数据:这个视图在排障时非常好用
SELECT * FROM iceberg_metadata('s3://my-lake/wh/ops/events');

这个组合我用得非常多:Spark 负责重型 ETL 和压实,DuckDB 负责临时分析和数据验证。同一份数据,两个引擎,零拷贝。这就是开放表格式的真正价值——存储和计算的解绑不再是 PPT 概念。


七、性能优化:从「能跑」到「跑得住」

7.1 压实策略:三个维度的取舍

压实(compaction)不是「跑一下 rewrite_data_files 就完事」。它有三个正交的目标,而且互相冲突:

目标手段副作用
消除小文件bin-pack 合并消耗 I/O,与写入争抢提交
消除删除向量重写带 DV 的文件写放大,且会重置文件级统计
优化数据聚集sort / zorder 重写成本最高,需要全局 shuffle

我的推荐是分层压实策略:

-- 第一层:高频(每小时),只做 bin-pack,只碰小文件
CALL lake.system.rewrite_data_files(
  table => 'ops.events',
  strategy => 'binpack',
  options => map(
    'min-input-files', '5',
    'target-file-size-bytes', '536870912',
    -- 只处理小于 target 的 25% 的文件,避免重写大文件
    'min-file-size-bytes', '134217728',
    -- 关键:部分进度提交,避免一次失败全部回滚
    'partial-progress.enabled', 'true',
    'partial-progress.max-commits', '10',
    'max-concurrent-file-group-rewrites', '10',
    -- 只压最近 2 天,冷数据别碰
    'where', "occurred_at >= current_date - interval 2 days"
  )
);

-- 第二层:中频(每天),清掉 DV 较多的文件
CALL lake.system.rewrite_data_files(
  table => 'ops.events',
  strategy => 'binpack',
  options => map(
    -- 一个文件被 >= 2 个删除文件引用就重写它
    'delete-file-threshold', '2',
    -- 或者按删除比例:删了 30% 以上就重写
    'delete-ratio-threshold', '0.3',
    'partial-progress.enabled', 'true'
  )
);

-- 第三层:低频(每周),做排序重写,优化查询局部性
CALL lake.system.rewrite_data_files(
  table => 'ops.events',
  strategy => 'sort',
  sort_order => 'user_id ASC NULLS LAST, occurred_at DESC',
  options => map(
    'rewrite-all', 'false',
    'partial-progress.enabled', 'true',
    'max-concurrent-file-group-rewrites', '5',
    'where', "occurred_at BETWEEN current_date - interval 8 days AND current_date - interval 1 day"
  )
);

partial-progress.enabled 这个参数被严重低估。默认情况下压实是一个大事务,重写 5TB 数据跑 3 小时,第 179 分钟失败 → 全部白干,而且这 3 小时里所有写入都在和你冲突。开启部分进度后,它会分批提交,失败只丢最后一批。大表压实必开。

7.2 元数据层的优化:别让 manifest 拖死规划

大表的查询慢,有时和数据一点关系没有,纯粹是元数据太碎。

-- 症状诊断:看 manifest 数量和平均大小
SELECT
  count(*) AS manifest_count,
  sum(added_data_files_count + existing_data_files_count) AS total_files,
  avg(length) / 1024 / 1024 AS avg_manifest_mb
FROM lake.ops.events.manifests;

-- 如果 manifest_count 上千而 avg_manifest_mb 只有几百 KB → 该重写了
CALL lake.system.rewrite_manifests(
  table => 'ops.events',
  -- 默认会用 spark 缓存,大表可能 OOM,关掉
  use_caching => false
);

判断标准:单个 manifest 目标大小 8MB 左右(commit.manifest.target-size-bytes 默认值),一个 manifest 大约管 8000–10000 个文件条目。如果你的表有 10 万个文件但有 2000 个 manifest,说明每个 manifest 只管 50 个文件——规划时要发 2000 次读请求,纯浪费。

配套的自动合并配置:

ALTER TABLE lake.ops.events SET TBLPROPERTIES (
  'commit.manifest-merge.enabled' = 'true',
  'commit.manifest.target-size-bytes' = '8388608',
  -- 提交时如果小 manifest 超过这个数就自动合并
  'commit.manifest.min-count-to-merge' = '100'
);

注意一个反直觉的坑:流式写入场景(Flink 每分钟提交)下,如果开启了 manifest 自动合并,每次提交都可能触发一次合并,导致提交延迟飙升。流式表建议关闭自动合并,改用定时的 rewrite_manifests。

7.3 快照过期与孤儿文件:省钱的部分

Iceberg 的快照永不自动删除。跑一年的表可能有几十万个快照,metadata.json 膨胀到几百 MB——每次提交都要读写这个文件,于是提交延迟从 100ms 涨到 30s。

-- 过期旧快照(保留 7 天 + 至少 100 个)
CALL lake.system.expire_snapshots(
  table => 'ops.events',
  older_than => TIMESTAMP '2026-08-06 00:00:00',
  retain_last => 100,
  -- 并行删除,大表必须给
  max_concurrent_deletes => 20
);

-- 清理孤儿文件(写失败残留的、没有任何元数据引用的文件)
-- ⚠️ 危险操作,见下方警告
CALL lake.system.remove_orphan_files(
  table => 'ops.events',
  older_than => TIMESTAMP '2026-08-10 00:00:00',
  dry_run => true   -- 先 dry_run 看清单!
);

remove_orphan_files 的致命陷阱:它的判断逻辑是「列举存储目录 vs 元数据引用」。如果此刻有一个正在运行的写入任务,它写完了数据文件但还没提交元数据,这些文件在 remove_orphan_files 看来就是孤儿——会被删掉,然后那个写入任务提交后表就指向了不存在的文件,表直接损坏。

铁律:older_than 至少设成最长写入任务耗时的 2 倍,并且永远先 dry_run=true。我见过因为这个参数设成 1 小时而搞挂生产表的事故。

7.4 对象存储层面的优化

湖仓性能的天花板往往在 S3 而不是计算。

ALTER TABLE lake.ops.events SET TBLPROPERTIES (
  -- 在文件路径前加哈希前缀,打散 S3 分区,避免请求限流
  'write.object-storage.enabled' = 'true',
  'write.data.path' = 's3://my-lake/data/events',

  -- 元数据文件数量上限
  'write.metadata.previous-versions-max' = '20',
  'write.metadata.delete-after-commit.enabled' = 'true',

  -- Parquet 行组大小:影响谓词下推粒度和内存
  'write.parquet.row-group-size-bytes' = '134217728',
  -- 页大小:小一点利于点查,大一点利于扫描
  'write.parquet.page-size-bytes' = '1048576',
  -- 字典编码阈值
  'write.parquet.dict-size-bytes' = '2097152',

  -- 只对高频过滤列建 min-max,全列建统计会让 manifest 膨胀
  'write.metadata.metrics.default' = 'counts',
  'write.metadata.metrics.column.user_id' = 'full',
  'write.metadata.metrics.column.occurred_at' = 'full',
  'write.metadata.metrics.column.event_type' = 'truncate(16)'
);

write.metadata.metrics.default = 'counts' 这一条能救命。 Iceberg 默认对前 100 列做 truncate(16) 的 min-max 统计。如果你的表有 500 列(宽表很常见),每个文件条目要存几百个上下界值,manifest 会膨胀到原来的 5–10 倍,规划直接变慢。只给真正用于过滤的列开 full,其他列只留 counts。

write.object-storage.enabled 也值得说:S3 的请求限流是按 key 前缀做的(每前缀约 3500 PUT / 5500 GET 每秒)。默认的 分区目录/文件名 布局会让同一分区的所有请求打在同一个前缀上。开启后 Iceberg 会插入一个哈希段,把请求打散到多个前缀。大规模写入场景这是必开项。

7.5 查询侧的调优

-- 让 split 大小匹配你的并行度
SET spark.sql.iceberg.split-size = 134217728;         -- 默认 128MB
SET spark.sql.iceberg.split-lookback = 10;
SET spark.sql.iceberg.split-open-file-cost = 4194304;

-- 向量化读(Parquet)
SET spark.sql.iceberg.vectorization.enabled = true;
SET spark.sql.iceberg.batch-size = 5000;

-- 规划并行度:大表规划时用多线程读 manifest
SET spark.sql.iceberg.planning.preserve-data-grouping = true;

split-open-file-cost 是一个被忽视的重要参数。它告诉规划器「打开一个文件的固定成本相当于读多少字节」。默认 4MB 意味着:如果你有大量小文件,规划器会把它们打包进同一个 split,而不是一个文件一个 task。小文件多的表把这个值调大到 8–16MB,能显著减少 task 数量和调度开销。


八、15 条生产踩坑清单

按「踩到的痛感」从高到低排:

1. remove_orphan_files 的 older_than 设太短,删掉了正在写入的文件。 后果是表损坏。铁律:至少设为最长写入任务耗时的 2 倍,永远先 dry_run。

2. V3 表没设 write.delete.granularity = 'file',DV 白开。 用了 partition 粒度就退化成了 V2 的行为,你以为升级了其实没有。升级后一定要用 inspect.delete_files() 确认 referenced_data_file 字段有值。

3. Flink checkpoint 间隔设成 10 秒。 一天 8640 个快照,metadata.json 膨胀,提交延迟指数上升。60s 是底线。

4. MERGE INTO 的 ON 条件没带分区谓词。 导致全表扫描。这是 MERGE 慢的第一大原因,改一行 SQL 能快 50 倍。

5. MERGE INTO 源表没去重。 Iceberg 会抛「a single row from target matched multiple rows」。必须在源侧用 ROW_NUMBER() 去重。

6. 分区粒度过细。 hours(ts) 在数据量不大时会产生海量小分区和小文件。经验值:单分区数据量目标 1GB 以上,宁可先粗后细(分区演进是元数据操作,可以改)。

7. 宽表开了全列 min-max 统计。 manifest 膨胀 10 倍,规划变慢。用 write.metadata.metrics.default = 'counts' 加白名单。

8. 忘记配 write.metadata.delete-after-commit.enabled。 老的 metadata.json 无限累积,某个凌晨你会发现一个表的元数据目录有 40 万个文件。

9. 高并发写同一张表,靠调大 commit.retry.num-retries 硬扛。 治不好。正确做法是按分区拆分写入任务,或用 branch 隔离后再合并。

10. 压实任务不开 partial-progress.enabled。 大表压实跑 3 小时后失败,全部回滚。

11. 以为 timestamp_ns 能无损容纳所有 timestamp 值。 纳秒 long 只能表示约 ±292 年范围。历史数据里有 1900 年之前或 2300 年之后的时间戳(脏数据里很常见)会溢出。迁移前先跑一遍范围检查。

12. Variant 类型没做 shredding 就上高频查询。 每次查询都要走二进制解析,虽然比 JSON 文本快,但比不上物化列。把 top-N 高频字段 shred 出来。

13. 依赖 V2 changelog 做增量,被压实污染。 压实一次下游全量重算。V3 上迁移到 _last_updated_sequence_number 水位法。

14. 用 SELECT count(*) 验证 V3 的元数据快路径,但表还有 equality delete 残留。 有 equality delete 存在时,行数无法从元数据算准。先跑一次全量压实把 equality delete 清干净。

15. 跨引擎写同一张表,但引擎的 Iceberg 版本不一致。 低版本引擎读不懂 V3 的 DV 会直接报错或者——更糟——静默地把已删除的行读出来。升级 V3 前,把所有读写方的 Iceberg 版本盘一遍,写下来。


九、总结与展望

9.1 V3 到底改变了什么

把这篇文章压缩成三句话:

  1. 删除向量把 MoR 的读放大从「和删除文件数成正比」压到了「常数一次 Range GET」,同时让规划器第一次能看见有效行数。这是「湖上能做频繁更新」的技术前提。

  2. 行级血缘给了每一行一个跨物理重写稳定的身份,让增量计算从「对比文件集合」变成「读一个序列号列」。这是「湖上能做流处理和物化视图」的技术前提。

  3. Variant + 地理类型 + 纳秒时间 + 列默认值 把湖仓的类型系统补到了和成熟数据库同一档位。这是「湖能当主库用」的技术前提。

三条合起来,指向同一个结论:湖仓不再是「数仓的廉价补充」,它开始具备成为主存储的资格。

9.2 还没解决的问题

诚实地说,V3 不是终点。留下的最大技术债是提交仍然是文件级的:

  • 每次提交都要写一个新的完整 metadata.json。表历史越长,这个文件越大,提交越慢
  • 每次提交至少要写 manifest list + manifest + metadata 三个文件,加上 Catalog 的 CAS,最少 4 次往返
  • 这决定了 Iceberg 的提交延迟下限在百毫秒到秒级,做不了真正的低延迟写入

社区关于 V4 的讨论方向也集中在这里:单文件提交(把多层元数据合并成一个可追加的结构)、自适应元数据(小表用简单结构,大表才分层)、更彻底的增量元数据。

同时值得关注的是表格式之外的竞争:DuckDB 团队的 DuckLake 提出了一个激进的思路——把元数据整个扔进一个 SQL 数据库,理由是「我们花了十年把数据库的能力重新实现在文件系统上,为什么不直接用数据库」。这个论点很难反驳。Iceberg 的分层文件元数据本质上是在「无事务的对象存储」这个约束下的最优解,而这个约束正在松动(S3 有了条件写入)。

9.3 给不同角色的行动建议

如果你在选型: V3 已经足够成熟,主流引擎(Spark、Flink、Trino、DuckDB、Snowflake)都有支持。选 Iceberg 的最大理由不是它技术最优,而是生态最广,不会把你锁死。

如果你已经在用 V2: 不要急着全表升到 V3。正确路径是:先在一张高更新率的表上试点 → 确认所有读写方版本支持 → 跑一次全量压实清掉 equality delete → 再改 format-version → 用 inspect.delete_files() 验证 DV 生效。

如果你在做增量管道: 优先级最高的动作是把 changelog 方案换成 row lineage 水位法。这个改造收益最大、风险最小。

如果你在被压实和小文件折磨: 先别急着优化压实,回头看两个东西——分区粒度是不是太细、Flink checkpoint 是不是太频繁。八成问题在这两个地方,压实只是在给上游的错误擦屁股。


最后说一句我的真实感受。 湖仓这个方向做了快十年,前七年在解决「怎么让文件系统上的表可靠」,后三年在解决「怎么让它快」。V3 是一个分界点——它标志着表格式这一层的核心问题基本收敛,接下来的竞争会转移到 Catalog 层(权限、多引擎协调、物化视图管理)和 Compute 层。

对写代码的人来说,这意味着一件很实在的事:你终于可以不用为了「行级更新」这一个需求,在数据栈里额外维护一套 OLTP 数据库了。 这笔账,省下来的运维时间比任何 benchmark 数字都值钱。

推荐文章

程序员茄子在线接单