编程 Apache Iceberg v3 深度拆解:当删除向量把行级更新从「小文件雪崩」拉回一次位运算——从 Puffin 二进制布局、Row Lineage 继承链到 Variant/地理类型与生产迁移全链路实战

2026-08-19 05:21:04 +0800 CST views 7

Apache Iceberg v3 深度拆解:当删除向量把行级更新从「小文件雪崩」拉回一次位运算

本文所有格式细节均以 Apache Iceberg 官方 Table Spec / Puffin Spec 为准(撰稿时 Java 实现最新发布版为 1.11.0)。文中标注「示意」的代码是为了讲清结构,标注「可运行」的代码我在本地跑过。

一、背景:为什么表格式非要出到 v3

先把版本语义摆清楚,很多人到今天还在混。

Iceberg 的 format-version 不是市场部起的名字,它是前向兼容的断点:版本号只在「老读端会读错新表」的时候才递增。官方 spec 对四个版本的定位是这样的:

版本官方定位核心能力状态
v1Analytic Data Tables用不可变文件(Parquet/Avro/ORC)管理大分析表,跟踪文件而非目录完成
v2Row-level Deletes引入 delete files,实现不重写数据文件的行级删除/更新完成
v3Extended Types and Capabilities新类型、默认值、多参数变换、Row Lineage、二进制删除向量、表加密密钥完成,社区已采纳
v4Metadata Structure and Representation重构元数据结构(含元数据字段的相对路径)在研,尚未正式采纳

注意最后一行:网上不少文章已经在「深度解析 v4」了,但 spec 页面上明明白白写着 under active development and has not been formally adopted。你现在能上生产的天花板是 v3。

v2 那笔技术债长什么样

v2 给了我们 merge-on-read(MoR):删一行不用重写 512MB 的 Parquet,只写一个「第几行被删了」的小文件。听起来很美,工程上留下三个病灶:

病灶一:小文件雪崩。 每次 commit 都可能产生新的 position delete 文件。CDC 场景下上游 MySQL 一分钟几千个 UPDATE,一天下来一张表挂几万个删除小文件是常态。对象存储上每个小文件都是一次 GET + 一次元数据往返。

病灶二:读端做 join。 v2 的 position delete 文件是「(file_path, pos)」的有序记录集,理论上一个删除文件可以引用任意多个数据文件。于是读端扫一个数据文件时,必须先把所有可能覆盖它的 delete 文件读进来、按 file_path 过滤、按 pos 归并排序,再和数据流做一次反连接。这个 join 的成本随删除文件数量线性上涨。

病灶三:写越勤,读越慢。 上面两条叠加的结果是一个反直觉的曲线:越是需要实时分析的表,查询性能越差。你被迫把 compaction 频率拉到极高,用计算成本换查询延迟。

v3 的删除向量(Deletion Vectors,下文简称 DV)就是冲着这三条来的。它的思路极其朴素:每个数据文件最多一个位图,位图第 P 位为 1 就表示第 P 行没了。 读端不再做 join,只做一次 bitmap.contains(pos) 的位运算。

二、核心概念:v3 到底加了哪六件事

按 spec 原文,v3 扩展了六个方向:

  1. 新数据类型:纳秒精度 timestamp_ns / timestamptz_nsunknownvariantgeometrygeography
  2. 列默认值initial-default / write-default
  3. 多参数变换:分区与排序支持多列输入的 transform(JSON 里从 source-id 变成 source-ids
  4. Row Lineage:行级血缘跟踪,_row_id + _last_updated_sequence_number
  5. 二进制删除向量:基于 Puffin 的 deletion-vector-v1
  6. 表加密密钥:table metadata 里的 encryption-keys 列表

这六件事里,前三件是「表达能力」,后三件是「工程能力」。真正改变架构的是 4 和 5,所以本文把主要篇幅给它们。

三、架构拆解 A:删除向量的字节级真相

3.1 位图为什么要分片

DV 支持 正的 64 位行位置(最高位必须为 0),但绝大多数文件行数远小于 2^32。如果直接上 64 位位图,稀疏场景下内存和序列化都是浪费。

spec 的做法是把 64 位位置劈成两半:

position (64bit, MSB must be 0)
┌──────────────┬──────────────┐
│  高 4 字节    │   低 4 字节   │
│    key       │  sub-position│
└──────────────┴──────────────┘
      │                │
      │                └─→ 存进该 key 对应的 32-bit Roaring Bitmap
      └─→ 定位是哪一个 Roaring Bitmap

查一个位置是否被删的逻辑就三步:取高 4 字节找 bitmap → 找不到就是没删 → 找到了则用低 4 字节测试是否 in-set。

这个设计非常务实:单文件行数 < 4.29 亿(2^32)时,整张 DV 只有一个 key=0 的 Roaring Bitmap,序列化后往往只有几十字节到几 KB;而当你真有超大文件时,它也不会崩,只是多几个分片。

3.2 Puffin:Iceberg 的「杂物间」文件格式

DV 不是独立文件格式,它寄生在 Puffin 里。Puffin 是 Iceberg 用来放「塞不进 manifest 的索引和统计信息」的通用容器,结构如下:

Magic  Blob₁  Blob₂ ... Blobₙ  Footer

Magic  = 0x50 0x46 0x41 0x31   ("PFA1", Puffin Fratercula arctica v1)
Footer = Magic  FooterPayload  FooterPayloadSize  Flags  Magic

要点:

  • FooterPayload 是 UTF-8 JSON(可选 LZ4 单帧压缩),描述所有 blob 的类型、偏移、长度、属性
  • FooterPayloadSize 是 4 字节小端有符号整数,压缩后长度
  • Flags 4 字节:byte 0 的最低位表示 FooterPayload 是否压缩,其余位保留为 0
  • 结尾又一次 Magic,让读端可以从文件尾巴反向定位 footer

所以读一个 Puffin 文件的正确顺序是:读最后 4 字节确认 Magic → 往前 4 字节读 Flags → 再往前 4 字节读 FooterPayloadSize → 定位 JSON → 解析出 blob 列表。一次尾部随机读就能拿到全部索引,这跟 Parquet 的 footer 思路完全一致。

3.3 deletion-vector-v1 blob 的字节布局

这是最容易被中文资料写错的地方,直接抄 spec:

┌────────────────────────────────────────────────────┐
│ 4 bytes  BIG-endian: len(magic + vector)           │
├────────────────────────────────────────────────────┤
│ 4 bytes  magic: D1 D3 39 64                        │
├────────────────────────────────────────────────────┤
│ N bytes  vector: Roaring "portable" 格式            │
│   ├ 8 bytes  LITTLE-endian: 32 位 bitmap 的个数     │
│   └ 每个 bitmap:                                    │
│       ├ 4 bytes LITTLE-endian: key                 │
│       └ 32-bit Roaring Bitmap 序列化字节            │
├────────────────────────────────────────────────────┤
│ 4 bytes  BIG-endian: CRC-32(magic + vector)        │
└────────────────────────────────────────────────────┘

三个坑:

  1. 长度字段和 CRC 是大端,Roaring 内部是小端。 spec 特意解释了原因:大端是为了兼容 Delta 表里已有的删除向量。这是一处有意为之的跨生态妥协。
  2. CRC 覆盖的是 magic + vector,不包含长度字段本身。
  3. blob 属性有强约束:必须包含 referenced-data-file(指向它作用的数据文件位置,且必须等于 table metadata 里记录的位置)和 cardinality(被删行数);必须省略 compression-codec,因为 deletion-vector-v1 不压缩。另外,Puffin 文件创建时 snapshot id 和 sequence number 还不知道,所以 blob metadata 里的 snapshot-idsequence-number 在 Puffin v1 必须写 -1

最后这一条特别值得玩味:它暴露了 Iceberg 一贯的设计哲学——先写数据,后定身份。因为乐观并发提交可能重试,任何在 commit 前就写死的 ID 都会导致重试时重写文件。所以能延后的全部延后。这个思路在 Row Lineage 里被发挥到极致。

3.4 manifest 侧的三个新字段

DV 在 delete manifest 里被单独跟踪,v3 为此新增:

字段含义
referenced_data_file该 DV 作用的数据文件位置(DV 必填;v2 那种只删一个文件的 position delete 也可选填)
content_offsetDV blob 在 Puffin 文件里的起始偏移,必须与 blob 实际偏移一致
content_size_in_bytesDV blob 长度

有了这三个字段,读端在 scan planning 阶段就知道「我要读的这个数据文件,对应的 DV 在哪个 Puffin 的哪个字节区间」,直接发一次 range GET 即可,不需要解析整个 Puffin footer 之外的东西。多个 DV 可以共存于同一个 Puffin 文件,且对它们引用哪些数据文件没有限制。

3.5 写端的硬约束(迁移时最容易踩)

spec 对 v3 写端下了几条死命令:

  • 每个数据文件最多一个 DV。 写端必须把新删除和旧 DV(以及旧的 position delete 文件)同步合并,不能像 v2 那样无脑追加。
  • 读端可以放心忽略 position delete 文件,只要该数据文件已经有 DV。
  • v3 表禁止新增 position delete 文件。
  • 从 v2 升级来的表,已有的 position delete 文件仍然有效;但当某数据文件首次生成 DV 时,必须把这些旧删除合并进 DV。
  • 那种「一个 position delete 文件覆盖多个数据文件」的老文件,必须一直保留在表元数据里,直到其中所有删除都被 DV 替代。

这几条翻译成运维语言就一句话:升级到 v3 之后的第一次 compaction,是一次不可逃避的、把历史删除债务一次性结清的动作。 你得给它留足资源和时间窗口。

3.6 可运行:手写一个 Puffin DV 解析器

调试线上表的时候,最想要的能力是「把这个 Puffin 文件里到底删了哪些行给我打出来」。按 spec 手写一个,不到 80 行:

# puffin_dv.py  —  可运行;依赖: pip install pyroaring
import json, struct, zlib, sys
from pyroaring import BitMap

MAGIC = b"\x50\x46\x41\x31"          # "PFA1"
DV_MAGIC = b"\xD1\xD3\x39\x64"

def read_footer(buf: bytes) -> dict:
    assert buf[:4] == MAGIC, "not a puffin file (head magic)"
    assert buf[-4:] == MAGIC, "not a puffin file (tail magic)"
    flags = buf[-8:-4]
    compressed = bool(flags[0] & 0x01)
    (size,) = struct.unpack("<i", buf[-12:-8])       # 小端有符号 4 字节
    payload = buf[-12 - size : -12]
    if compressed:
        import lz4.frame                              # pip install lz4
        payload = lz4.frame.decompress(payload)
    return json.loads(payload.decode("utf-8"))

def parse_dv_blob(blob: bytes) -> BitMap:
    (length,) = struct.unpack(">i", blob[:4])        # 大端:magic+vector 长度
    body = blob[4 : 4 + length]
    (crc,) = struct.unpack(">I", blob[4 + length : 8 + length])
    assert zlib.crc32(body) & 0xFFFFFFFF == crc, "DV CRC mismatch"
    assert body[:4] == DV_MAGIC, "bad DV magic"
    vec = body[4:]
    (n,) = struct.unpack("<q", vec[:8])              # 小端:bitmap 个数
    out, off = BitMap(), 8
    for _ in range(n):
        (key,) = struct.unpack("<I", vec[off : off + 4]); off += 4
        bm = BitMap.deserialize(vec[off:])            # 32-bit portable roaring
        off += len(bm.serialize())
        for p in bm:
            out.add((key << 32) | p)                  # 还原 64 位位置
    return out

if __name__ == "__main__":
    raw = open(sys.argv[1], "rb").read()
    meta = read_footer(raw)
    for b in meta["blobs"]:
        if b["type"] != "deletion-vector-v1":
            continue
        props = b.get("properties", {})
        blob = raw[b["offset"] : b["offset"] + b["length"]]
        bm = parse_dv_blob(blob)
        print(f"[DV] file={props.get('referenced-data-file')}")
        print(f"     cardinality(prop)={props.get('cardinality')} parsed={len(bm)}")
        print(f"     deleted positions (first 20)={list(bm)[:20]}")

跑起来长这样:

$ python puffin_dv.py 00000-3-9c9e...-deletes.puffin
[DV] file=s3://lake/db/orders/data/00042-...-a1b2.parquet
     cardinality(prop)=1731 parsed=1731
     deleted positions (first 20)=[3, 17, 18, 19, 40, 41, ...]

cardinality 属性和实际解析出的基数不一致,就说明写端有 bug 或者文件被截断了——这是一条很实用的线上体检项。

四、架构拆解 B:Row Lineage 的四级继承链

如果说 DV 是性能story,Row Lineage 是语义 story。它回答一个此前湖仓答不上来的问题:这一行是什么时候进表的、上一次被谁改的?

4.1 两个保留列

v3 及以后,Iceberg 表必须为所有新建行跟踪两个字段(引擎侧必须同时维护表级 next-row-id):

  • _row_id:表内唯一的 long 标识,行首次写入时通过继承赋值
  • _last_updated_sequence_number:最后一次修改该行的那次 commit 的 sequence number,首次写入或修改时通过继承赋值

4.2 为什么是继承,不是直接写死

spec 给的理由值得每个做分布式存储的人抄进笔记本:

因为 commit sequence number 和起始 row ID 在快照成功提交之前都还没确定。用继承可以在这些值未知时先把数据文件和 manifest 写出去,这样乐观提交重试时不必重写文件

这就是前面 Puffin blob 里 snapshot-id = -1 的同一套哲学。展开成一条完整的赋值链:

table metadata: next-row-id  = 1000
        │  (提交时把当前 next-row-id 交给快照)
        ▼
snapshot: "first-row-id": 1000
        │  (写 manifest list 时,按前序新 manifest 的
        │    added_rows_count + existing_rows_count 累加分配)
        ▼
manifest: first_row_id = 1000 / 1125 / 1225 ...
        │  (按数据文件在 manifest 中的顺序继续累加)
        ▼
data file: first_row_id
        │
        ▼
row: _row_id = data_file.first_row_id + _pos

spec 里给的例子(表初始 next-row-id = 1000):

manifestadded_rows_countexisting_rows_countfirst_row_id
existing750925
added1100251000
added201001125
added3125251225

added1 拿到和快照相同的 1000;added2 的起点 = 1000 + (100 + 25) = 1125;added3 = 1125 + (0 + 100) = 1225。已存在的 manifest 保留它被加入表时分配的值(925),不会被重排。

读端规则同样简洁:

  • 数据文件里若 _row_id 为 null(或干脆没有这一列),读端按 first_row_id + _pos 现场算
  • _last_updated_sequence_number 为 null,读端取该数据文件 manifest entry 的 sequence_number
  • 只有新行的数据文件可以完全省略这两列,读端应当视其存在且全为 null

也就是说,append-only 的写入路径几乎零额外开销:不写这两列,读时算出来就行。

4.3 行搬家时的三条规则

当一行因为 compaction、分区演进等原因被搬到新数据文件,写端必须:

  1. 把该行非 null 的 _row_id 原样复制到新文件
  2. 如果这次写入修改了该行,把 _last_updated_sequence_number 置为 null(好让本次修改的 sequence number 通过继承生效)
  3. 如果没修改,把原来那个非 null 的 _last_updated_sequence_number 原样复制过去

这三条保证了「compaction 不改变行的身份,只改变它的物理位置」。这正是构建增量物化视图的基石。

4.4 一个重要的例外:equality delete 不参与血缘

spec 明确:equality delete 造成的更新不跟踪血缘。原因是使用 equality delete 的引擎(典型是 Flink 上游 CDC 写入)刻意不去读现存数据就直接写变更,因此拿不到原行的 row ID。这类更新在语义上被当作「旧行整体删除 + 全新行插入」。

这条限制的现实含义很硬核:如果你的链路重度依赖 equality delete(比如 Flink upsert sink),你就别指望 Row Lineage 能给你一条连续的行历史。 想要血缘,就得改成 read-then-write 的 MoR/CoW 路径。这是选型时必须提前拍板的取舍,不是上线后能调参解决的。

4.5 Row Lineage 能换来什么

  • 真增量 CDC:对比两个快照,用 _last_updated_sequence_number > S 直接筛出变更行,不再依赖上游打时间戳字段
  • 审计与合规:GDPR 删除请求的执行证据链
  • 增量物化视图刷新:只重算受影响的行,而不是分区级重算
  • 数据质量归因:某条脏数据是哪次 commit 引入的,一查就知道
-- 示意:拿到自 sequence number 8421 以来发生过变化的行
SELECT _row_id, _last_updated_sequence_number, order_id, status
FROM   lake.db.orders
WHERE  _last_updated_sequence_number > 8421;

-- 配合元数据表定位那次提交
SELECT snapshot_id, sequence_number, operation, committed_at, summary
FROM   lake.db.orders.snapshots
ORDER  BY sequence_number DESC LIMIT 20;

五、类型系统:把「半结构化」和「地理」收进标准

5.1 variant

variant 存半结构化数据,同一列在不同行里的结构和类型可以不同。关键事实:

  • variant 的类型与二进制编码定义在 Parquet 项目里,Iceberg 复用它(目前支持 V1 编码)
  • 它既不是 primitive,也不是 nested type,是独立的第三类
  • 比 JSON 更宽:原生支持 date、timestamp、timestamptz、binary、decimal
  • 内部可嵌套 array(元素类型不固定)和 object(字段名可变、值类型任意)
  • spec 单独定义了「Bounds for Variant」,也就是说 variant 的子字段可以参与统计与裁剪

对比过去的做法(把 JSON 塞进 string 列)差别巨大:string 方案下引擎只能做全列扫描 + 运行时解析,而 variant 有编码规范和统计边界,可以做谓词下推

-- 示意:v3 表 + variant 列
CREATE TABLE lake.db.events (
  event_id   BIGINT,
  occurred_at TIMESTAMP,
  payload    VARIANT           -- 结构随业务演进,不用 ALTER TABLE
) USING iceberg
PARTITIONED BY (days(occurred_at))
TBLPROPERTIES (
  'format-version' = '3',
  'write.delete.mode' = 'merge-on-read',
  'write.update.mode' = 'merge-on-read',
  'write.merge.mode'  = 'merge-on-read'
);

注意:variant 的 SQL 语法与函数在各引擎间进度不一,落地前务必查你那套引擎的版本支持矩阵(官方有 Implementation Status 页)。

5.2 geometry / geography

v3 把地理类型收编进 spec,并且带了两个专门的概念:

  • CRS(坐标参考系)
  • Edge-Interpolation Algorithm(边插值算法)——geography 类型需要它来定义两点之间的边如何在球面上插值

spec 甚至给了 Appendix G: Geospatial Notes,还专门定义了「Bounds for Geometry and Geography」。这意味着地理查询的分区裁剪和文件裁剪是标准的一部分,而不是各家引擎的私货。以前做 GIS 的同学在湖上只能用 WKB + string,现在有正经类型了。

5.3 默认值:加列不重写数据

v3 允许 struct 字段(含顶层 schema)带两个默认值:

  • initial-default:用于填充加列之前已写入的所有记录
  • write-default:用于填充加列之后写入但写端没给值的记录

initial-default 只在往已有 schema 加列时设置;write-default 初值等于它,之后可以通过 schema evolution 改。若可选字段两个默认值都没设,则默认为 null(向后兼容)。

最关键一句:这套机制产生标准 SQL 的默认值行为,而不需要重写数据文件。 以前给 PB 级表加一个「NOT NULL DEFAULT 0」的列,等于一次全表重写;现在只改元数据。

另外 spec 明确:写数据文件时不允许省略已知字段,默认值的变更只影响未来的记录。别把 write-default 当成「稀疏列压缩」用。

5.4 多参数变换

分区/排序的 transform 从单列走向多列,JSON 序列化上体现为:多参数变换必须写 source-ids,单参数变换写 source-id

顺便复习一下 v3 的 transform 全集(含纳秒类型的加入):

transform说明关键点
identity原值除 geometry/geography 外的任意 primitive
bucket[N]hash mod N32 位 Murmur3 x86 变体,seed = 0;先丢符号位保证非负
truncate[W]截断到宽度 Wint/long/decimal/string/binary
year/month/day/hour时间抽取已扩展到 timestamp_ns / timestamptz_ns
void恒为 null用于 v1 表里「逻辑删除」某个分区字段

day 的结果类型是 date,但 spec 要求读端也必须接受 int(按 1970-01-01 起的天数解释)——这类兼容细节是自己实现读端时最容易翻车的地方。

六、代码实战:从建表到看见 DV

下面这套流程我按「能复现」的顺序排,用 Spark + PyIceberg。

6.1 起一个 REST catalog 的 Spark 环境

# 示意:本地 minio + rest catalog
export ICEBERG_VERSION=1.11.0
spark-sql \
  --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:${ICEBERG_VERSION} \
  --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://lake/ \
  --conf spark.sql.catalog.lake.io-impl=org.apache.iceberg.aws.s3.S3FileIO

6.2 建 v3 表并制造行级删除

CREATE TABLE lake.db.orders (
  order_id   BIGINT,
  user_id    BIGINT,
  amount     DECIMAL(18,2),
  status     STRING,
  created_at TIMESTAMP
) USING iceberg
PARTITIONED BY (days(created_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',
  'write.target-file-size-bytes' = '536870912',
  'write.distribution-mode'   = 'hash'
);

INSERT INTO lake.db.orders VALUES
 (1, 1001, 99.00, 'PAID',   TIMESTAMP '2026-08-18 10:00:00'),
 (2, 1002, 12.50, 'CREATED',TIMESTAMP '2026-08-18 10:05:00'),
 (3, 1001, 88.80, 'CREATED',TIMESTAMP '2026-08-18 11:00:00');

-- 触发 DV 生成
DELETE FROM lake.db.orders WHERE status = 'CREATED' AND amount < 20;
UPDATE lake.db.orders SET status = 'SHIPPED' WHERE order_id = 1;

6.3 用元数据表验证 DV 落地

-- 看删除文件的形态:content=1 是 position delete,DV 在 v3 里也走 delete manifest
SELECT file_path, file_format, content, record_count,
       referenced_data_file, content_offset, content_size_in_bytes
FROM   lake.db.orders.all_delete_files;

-- 看每个数据文件当前挂了几个删除
SELECT file_path, record_count, file_size_in_bytes
FROM   lake.db.orders.files;

-- 快照与 sequence number(Row Lineage 的时间轴)
SELECT snapshot_id, sequence_number, operation, summary['added-delete-files']
FROM   lake.db.orders.snapshots ORDER BY sequence_number;

referenced_data_file / content_offset / content_size_in_bytes 三列非空,就说明 DV 路径真的生效了。把 file_path 拿去喂第 3.6 节那个解析器,就能看到具体被删的行位置。

6.4 PyIceberg 侧:读元数据、算删除率

# 可运行;依赖: pip install "pyiceberg[s3fs,pyarrow]"
from pyiceberg.catalog import load_catalog

cat = load_catalog("lake", **{
    "type": "rest",
    "uri": "http://localhost:8181",
    "warehouse": "s3://lake/",
})
tbl = cat.load_table("db.orders")

print("format-version:", tbl.metadata.format_version)
print("current snapshot:", tbl.metadata.current_snapshot_id)

# 扫描计划:拿到每个数据文件挂了哪些 delete 文件
total_rows = total_deleted = 0
for task in tbl.scan().plan_files():
    dels = list(task.delete_files)
    total_rows += task.file.record_count
    total_deleted += sum(getattr(d, "record_count", 0) or 0 for d in dels)
    if dels:
        print(f"{task.file.file_path.split('/')[-1]}"
              f"  rows={task.file.record_count}  deletes={len(dels)}")

ratio = total_deleted / max(total_rows, 1)
print(f"\n删除率 = {ratio:.2%}  (>10% 建议触发 compaction)")

这个「删除率」指标非常值得做成日常巡检项。DV 让读放大变成了位运算,但被逻辑删除的行仍然占着物理空间和 I/O 带宽——DV 解决的是「合并删除的 CPU 成本」,没解决「白读一遍再丢掉」的带宽成本。

6.5 Java 侧的 RowDelta(示意)

// 示意:自己写入 DV 时的提交骨架
Table table = catalog.loadTable(TableIdentifier.of("db", "orders"));

DeleteFile dv = FileMetadata.deleteFileBuilder(table.spec())
        .ofPositionDeletes()
        .withPath("s3://lake/db/orders/data/00042-deletes.puffin")
        .withFileSizeInBytes(size)
        .withReferencedDataFile(dataFilePath)   // v3 必填
        .withContentOffset(offset)              // v3 必填
        .withContentSizeInBytes(blobLength)     // v3 必填
        .withRecordCount(cardinality)
        .withFormat(FileFormat.PUFFIN)
        .build();

table.newRowDelta()
     .addDeletes(dv)
     .validateFromSnapshot(baseSnapshotId)      // 乐观并发校验
     .commit();

真实项目里这块基本由引擎接管,但当你要写自定义摄入器(比如自研 CDC sink)时,这三个 v3 字段漏一个,读端就会静默读到脏数据——因为老读端会把它当成一个「覆盖多文件的 position delete」去处理。

七、性能优化:一份可直接抄的调优清单

7.1 CoW 还是 MoR:先按写放大定,再按查询延迟修

场景推荐理由
天级批量 overwritecopy-on-write反正整分区重写,CoW 读最快
高频小批量 UPDATE / DELETE(CDC)merge-on-read避免每次删一行重写 512MB
合规删除(低频、随机)merge-on-read + 定期 compactionDV 让随机删几乎零成本
查询延迟敏感、写入极低copy-on-write没有删除文件,无需合并
ALTER TABLE lake.db.orders SET TBLPROPERTIES (
  'write.delete.mode' = 'merge-on-read',
  'write.update.mode' = 'merge-on-read',
  'write.merge.mode'  = 'merge-on-read',
  -- 删除粒度:file 让 DV 只针对单个数据文件,利于精准合并
  'write.delete.granularity' = 'file'
);

7.2 compaction:不要等到删除率报警

-- 数据文件重写(带部分进度提交,避免一次性大事务失败全丢)
CALL lake.system.rewrite_data_files(
  table => 'db.orders',
  strategy => 'sort',
  sort_order => 'created_at DESC NULLS LAST',
  options => map(
    'target-file-size-bytes', '536870912',
    'min-input-files', '5',
    'delete-file-threshold', '4',        -- 挂了 4 个以上删除文件就重写
    'partial-progress.enabled', 'true',
    'partial-progress.max-commits', '10',
    'max-concurrent-file-group-rewrites', '8',
    'rewrite-job-order', 'bytes-desc'
  )
);

-- manifest 重写:清单文件太碎会拖慢 scan planning
CALL lake.system.rewrite_manifests('db.orders');

-- 快照过期:真正释放存储的一步
CALL lake.system.expire_snapshots(
  table => 'db.orders',
  older_than => TIMESTAMP '2026-08-12 00:00:00',
  retain_last => 10
);

-- 孤儿文件清理:务必先确认没有并发写入
CALL lake.system.remove_orphan_files(table => 'db.orders', dry_run => true);

三条经验:

  1. delete-file-threshold 是 MoR 表最重要的一个旋钮。 它决定「一个数据文件挂多少删除就值得重写」。DV 时代阈值可以比 v2 放宽,但别放成无穷。
  2. partial-progress.enabled 在大表上几乎必开。 否则跑 4 小时的 compaction 在最后一次提交冲突时全部回滚。
  3. expire_snapshots 才是省钱的那一步。 只跑 compaction 不过期快照,旧文件被历史快照引用着,存储只涨不降。

7.3 元数据层的隐性成本

ALTER TABLE lake.db.orders SET TBLPROPERTIES (
  'commit.manifest.target-size-bytes'  = '8388608',
  'commit.manifest.min-count-to-merge' = '100',
  'commit.manifest-merge.enabled'      = 'true',
  'write.metadata.delete-after-commit.enabled' = 'true',
  'write.metadata.previous-versions-max'       = '50',
  'history.expire.max-snapshot-age-ms'  = '432000000',  -- 5 天
  'history.expire.min-snapshots-to-keep' = '10'
);

高频提交的流式表最典型的病是 metadata.json 版本堆到几万个,每次 loadTable 都变慢。write.metadata.delete-after-commit.enabled + previous-versions-max 直接治这个。

7.4 读侧

  • read.split.target-size(默认 128MB):对象存储上适当放大到 256MB,减少任务数与请求数
  • 分区设计优先用 hidden partitioning + bucket,避免高基数列直接 identity 分区导致的分区爆炸
  • 排序写入(strategy => 'sort')能显著提升 min/max 裁剪效果,比堆机器便宜得多
  • Puffin 里还能放 apache-datasketches-theta-v1(NDV 草图),给 CBO 提供基数估计——这是很多人忽略的免费收益

7.5 一句反常识的话

DV 让「删」变便宜,但没让「查」变自由。 我见过团队升到 v3 后把 compaction 频率砍掉一半,结果查询变慢:因为删除率累积到 30% 时,你有三成的 I/O 是读进来就丢掉的。正确的做法是:把 compaction 触发条件从「删除文件个数」改成「删除率 + 文件大小分布」,而不是直接放松。

八、迁移到 v3 的检查清单

按我的经验,踩坑顺序基本固定:

  1. 先查读端。 所有会读这张表的引擎(Spark/Flink/Trino/StarRocks/Doris/DuckDB/自研服务)都必须支持 v3 的 DV 与 row lineage 字段。format-version 升级不可回退,这是单向门。以官方 Implementation Status 页为准,不要凭博客判断。
  2. 盘点 equality delete 依赖。 如果 Flink upsert 链路在用 equality delete,明确接受「这部分数据没有 row lineage」,或者改造写入路径。
  3. 给第一次 compaction 留窗口。 升级后旧的 position delete 文件需要被合并进 DV,尤其那些跨多数据文件的老删除文件会一直挂在元数据里直到被完全替代。
  4. 自研摄入器逐个补齐三个新字段referenced_data_file / content_offset / content_size_in_bytes),并确保「一个数据文件最多一个 DV」的同步合并逻辑。
  5. 写端必须维护 next-row-id 这不是可选项,v3 表的写端如果不维护它,行血缘会直接错乱。
  6. 加密密钥先想清楚 KMS。 encryption-keys 里存的是 base64 的加密后密钥元数据,格式由表的加密方案决定,可以是 KMS 特定的包装格式;encrypted-by-id 指向包装它的那把钥匙。密钥轮换策略要在上线前定,不要事后补。
  7. 灰度:先升一张非核心大表,跑满一个完整的「写入 → 删除 → compaction → 过期 → 全量查询校验」周期再推广。

九、几个高频疑问的正面回答

Q:DV 和 Delta Lake 的 Deletion Vectors 是一回事吗?
不是同一套实现,但刻意保持了互操作友好。最直接的证据就是 spec 里那句解释:长度与 CRC 字段用大端,是为了兼容 Delta 表里已有的删除向量。两边都用 Roaring Bitmap 表达「哪些位置被删」,差别在容器(Iceberg 走 Puffin,并在 delete manifest 里用三个新字段索引到 blob 的字节区间)。这意味着做跨格式转换的工具链,位图部分可以复用。

Q:v3 表还能有 equality delete 吗?
能。spec 禁止的是「v3 表新增 position delete 文件」,equality delete 并不在禁令里。但代价是这类更新不参与 Row Lineage,语义上被当作「删旧行 + 插新行」。所以你会遇到一个真实的架构分叉:要极致的流式写入吞吐(不读现存数据),还是要行级血缘。这两个目标在 v3 里是互斥的,不要指望调参解决。

Q:升级 format-version 会不会锁死?
会。升级是单向的,没有官方降级路径。真要回退只能重建表(CTAS 到一张 v2 表再切换名字),代价是全量重写 + 历史快照丢失。所以第 8 节的第 1 条「先查所有读端」不是形式主义。

Q:_row_id 能当业务主键用吗?
不能。它是表内唯一的物理身份标识,不承载业务语义;且 equality delete 路径下会「换号」。它适合做的是审计、增量 diff、行级归因,不适合被下游系统当外键存起来。

Q:DV 会不会导致「删了但空间没释放」?
一定会,而且这是设计使然。DV 只是逻辑标记,物理数据还在 Parquet 里。真正释放空间需要两步:compaction 把存活行重写成新文件,再 expire_snapshots 让旧文件不再被任何快照引用。很多团队只做第一步,然后困惑存储账单为什么不降。

Q:一个 Puffin 文件放多少个 DV 合适?
spec 允许多个 DV 共存于同一 Puffin,且不限制它们引用哪些数据文件。工程上的权衡是:放太少 → 小文件多;放太多 → 单个 Puffin 变大,虽然靠 content_offset + content_size_in_bytes 能精准 range GET,但写端每次合并 DV 时的重写成本会上升。按分区聚合、单文件控制在几 MB 量级是比较稳的起点。

十、总结与展望

把这篇的技术判断压缩成几条:

  • v3 是 Iceberg 从「能存」走向「能改」的分水岭。 v2 给了行级删除的语义,但工程代价高得让很多团队退回 CoW;v3 用 DV 把这个代价降到位图级别。
  • DV 的价值不在位图本身,而在「至多一个 DV / 数据文件」这条约束。 它把读端从「不定数量的 delete 文件 join」变成「一次 range GET + 一次 contains」,复杂度从 O(n) 掉到 O(1)。
  • Row Lineage 是被严重低估的特性。 它让湖仓第一次具备行级身份,增量视图、审计、CDC 下游都能从中受益。它的继承式赋值设计(延迟到 commit 才定 ID)是分布式乐观并发的教科书级案例。
  • 类型系统的扩展(variant / geospatial / 默认值)解决的是「不用为了迁就格式而扭曲建模」。 加列不重写数据这一条,对 PB 级表来说就是几十万块钱的差别。
  • v4 别急。 spec 明说在研未采纳,方向是元数据结构重构与相对路径(后者对表整体搬迁、跨环境克隆意义很大)。现在该做的是把 v3 用透。

最后说点行业层面的观察:v3 的很多设计(比如 DV 的大端长度字段刻意兼容 Delta)说明表格式的战争正在从「谁的格式赢」转向「怎么互操作」。对我们写代码的人来说,这是好事——未来的竞争在引擎和 catalog,而不是在文件布局上重复造轮子。

参考

  • Apache Iceberg Table Spec(Format Versioning / Row Lineage / Deletion Vectors / Partition Transforms / Appendix E)
  • Apache Iceberg Puffin Spec(File structure / Footer / deletion-vector-v1 blob type)
  • Apache Iceberg Releases 与 Implementation Status 页
  • Roaring Bitmap portable serialization format

推荐文章

mysql int bigint 自增索引范围
2024-11-18 07:29:12 +0800 CST
Golang 几种使用 Channel 的错误姿势
2024-11-19 01:42:18 +0800 CST
程序员茄子在线接单