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

2026-08-13 22:29:43 +0800 CST views 9

如果你在 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 INTOALTER 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_filesrewrite_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_idpayload.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 INTOALTER 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_filesolder_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 INTOON 条件没带分区谓词。 导致全表扫描。这是 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 数字都值钱。

推荐文章

Vue3 组件间通信的多种方式
2024-11-19 02:57:47 +0800 CST
Vue3 中提供了哪些新的指令
2024-11-19 01:48:20 +0800 CST
关于 `nohup` 和 `&` 的使用说明
2024-11-19 08:49:44 +0800 CST
JavaScript 实现访问本地文件夹
2024-11-18 23:12:47 +0800 CST
程序员茄子在线接单