如果你在 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 | 基数 > 4096 | 8KB 定长位图 | ~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 用「字典下标」引用,不重复存字符串
好处:
- 不需要 parse 文本。取
payload.user.id是二进制偏移跳转,不是字符串扫描 - key 字符串只存一次。埋点数据里
"event_timestamp"这种长 key 出现 100 万次,只在 metadata 里存一份 - 保留类型。JSON 里的
123和"123"是不同的,round-trip 不丢信息 - 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 同步、每分钟级 upsert | MoR | 写放大不可承受 |
| GDPR 删除(按用户删,极稀疏) | MoR + DV | DV 对稀疏删除是最优编码 |
| 按时间范围批量删(连续段) | MoR + DV | Run 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 |
| 能否文件级剪枝 | 不能 | 能(序列号在元数据里) |
6.6 Flink 流式 Upsert
// 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 的三个必踩坑,提前说:
checkpoint 间隔就是提交间隔。 设 10s 你一天会产生 8640 个快照和天量小文件。低于 60s 我基本不建议,除非你有配套的实时压实。
upsert-enabled=true要求表有主键,且必须配write.distribution-mode=hash。 否则同一个 key 的不同版本会落到不同的 writer subtask,导致删除标记和数据错位——这类数据不一致极难排查。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 到底改变了什么
把这篇文章压缩成三句话:
删除向量把 MoR 的读放大从「和删除文件数成正比」压到了「常数一次 Range GET」,同时让规划器第一次能看见有效行数。这是「湖上能做频繁更新」的技术前提。
行级血缘给了每一行一个跨物理重写稳定的身份,让增量计算从「对比文件集合」变成「读一个序列号列」。这是「湖上能做流处理和物化视图」的技术前提。
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 数字都值钱。