编程 Apache Iceberg V3 深度拆解:删除向量、行血缘与 Variant 类型,湖仓表格式的第三次范式转移

2026-08-02 01:39:04 +0800 CST views 6

Apache Iceberg V3 深度拆解:删除向量、行血缘与 Variant 类型,湖仓表格式的第三次范式转移

2026 年,「湖仓一体」已经不是 PPT 词汇,而是生产环境的默认配置。Iceberg 的 Java 实现走到了 1.11.0,Hive 4.2.0 宣布支持 V3 的 Deletion Vectors / 列默认值 / Variant,Snowflake、OceanBase 这类"重"系统也纷纷跟进 V3 格式。但绝大多数团队还停留在 V2,甚至还有人在用 V1。

这篇文章不讲"Iceberg 是什么"。我们直接钻进 V3 规范的骨头里:删除向量为什么能把 MOR 的读放大摁下去、行血缘凭什么可以不读数据就分配全局唯一 ID、Variant 二进制编码怎么解决半结构化数据的世纪难题。以及最重要的——你现在该不该升,怎么升,升完会踩哪些坑


一、先说痛点:V2 的 Merge-on-Read 是怎么把自己玩坏的

1.1 一个真实的性能坍塌现场

假设你有一张 CDC 同步过来的订单表,Iceberg V2 格式,Merge-on-Read(MOR)模式。数据量不大,2TB,1 万个 Parquet 数据文件。上游 MySQL binlog 每 5 分钟一批,通过 Flink 写进来。

刚上线的时候一切美好,查询 3 秒返回。两周之后,同样的查询变成了 90 秒。

你去看执行计划,发现 Scan 阶段读了 1 万个数据文件 + 8 万个 delete 文件

这就是 V2 position delete 的死穴。

1.2 V2 的删除模型:两种 delete file

Iceberg V2 引入了两种删除文件:

Position Delete(位置删除):记录 (file_path, pos) 二元组,表示"某个数据文件的第 N 行被删了"。

# position delete 文件的逻辑内容(Parquet 格式存储)
file_path                                    | pos
---------------------------------------------|------
s3://bucket/data/00000-1-abc.parquet         | 42
s3://bucket/data/00000-1-abc.parquet         | 1337
s3://bucket/data/00001-2-def.parquet         | 7

Equality Delete(等值删除):记录"某些列等于某些值的行被删了",典型是 id = 12345。Flink CDC 大量使用这种。

# equality delete 文件
order_id
---------
12345
12346

1.3 死穴在哪

问题不在于删除文件本身,而在于它们的数量是可以无限膨胀的,而读取时必须全部加载

看 V2 规范定义的读取语义:

对于一个数据文件 D,需要应用的 position delete 文件集合 = 所有满足以下条件的 delete file:

  • sequence_number >= D.sequence_number
  • 分区匹配
  • (可选)通过 referenced_data_file 或列统计过滤

关键词是集合。一个数据文件可能被 N 个 delete file 覆盖:第一批删了第 42 行,第二批删了第 1337 行,第三批又删了第 42 行(重复删除是合法的)……每次 commit 都可能产生新的 delete file。

于是读路径变成了:

读取 data file D
  -> 找到所有相关 delete files [d1, d2, ..., dN]
  -> 逐个打开、解码、合并成一个位置集合
  -> 归并排序 / 构建 HashSet
  -> 对 D 的每一行做 lookup 过滤

这里有三层放大:

  1. 文件数放大:N 个 delete file = N 次对象存储 GET。S3 单次 GET 的 P99 在 50-150ms,8 万个文件哪怕并发 200 也要几十秒。
  2. 内存放大:Parquet position delete 每行是 (string, long),1000 万条记录光 file_path 解码后的 Java String 就能吃掉几百 MB 堆内存。
  3. CPU 放大:合并 N 个有序流需要堆排序,或构建一个巨大的 HashSet。

更恶心的是:你无法在 planning 阶段准确知道要读多少 delete file,因为过滤条件只能靠列统计做粗粒度裁剪。

1.4 业界的临时解法,以及为什么都不够

  • 疯狂 compaction:每小时跑 rewrite_data_files 把 delete 合并回去。代价是巨大的写放大——为删 1000 行重写 100GB。
  • 改用 Copy-on-Write:写时直接重写整个数据文件。CDC 场景下每 5 分钟重写一批,成本爆炸。
  • 拄 Delta Lake 的 Deletion Vector:Databricks 在 Delta Lake 2.3(2023)就用 Roaring Bitmap 存删除位置。Iceberg 社区看了两年。

V3 的答案就是解法 C,但做得更彻底:不是"多了一种删除方式",而是"position delete 文件在 V3 里被彻底废弃,只剩 DV"


二、Iceberg 元数据结构速通(读懂后面必需)

在拆 V3 之前,必须把 Iceberg 的元数据分层刻在脑子里。很多人只知道"Iceberg 有快照",但不清楚每一层的物理形态。

Catalog(Hive Metastore / REST Catalog / Glue / JDBC / Nessie)
   |  只存一个指针:current metadata location
   v
v42.metadata.json                      <- 表元数据,JSON,全量重写
   |  schemas / partition-specs / sort-orders
   |  snapshots[] / refs{} / properties
   v
snap-{id}-{seq}-{uuid}.avro            <- manifest list,Avro
   |  每行 = 一个 manifest 文件的摘要
   |  含分区值域(partition summaries)用于裁剪
   v
{uuid}-m0.avro                         <- manifest 文件,Avro
   |  每行 = 一个 data file 或 delete file 的完整元数据
   |  含每列的 min/max/null_count/value_count
   v
00000-0-{uuid}.parquet                 <- 实际数据
00000-0-{uuid}.puffin                  <- V3 的删除向量(新增)

2.1 为什么是这个结构

核心目的是让 query planning 不需要 LIST 对象存储。Hive 表格式的原罪就是“目录即分区”,planning 时要递归 LIST S3,而 S3 的 LIST 是分页的、最终一致的、慢的。一张 10 万分区的 Hive 表,光 planning 就要几分钟。

Iceberg 把文件列表写进不可变的 Avro,planning 变成:读 JSON → 读 manifest list → 按分区摘要裁 manifest → 按列统计裁 data file。全程 O(相关文件数),且都是 GET 而非 LIST。

2.2 sequence number:并发控制的锚点

这是理解 delete 语义的关键。每次 commit 会产生一个单调递增的 sequence_number

  • 数据文件带着它被写入时的 sequence number
  • 删除文件也带着 sequence number
  • 规则:delete file 只能作用于 sequence_number 小于等于它自己的 data file

这条规则保证了正确性:新写入的数据不会被旧的删除操作误伤。想象两个并发事务,T1 删除 id=5,T2 插入 id=5。如果 T2 的 seq 更大,那么 T1 的 delete file 不会应用到 T2 的数据文件上,id=5 存活——这符合"删除发生在插入之前"的时间语义。

2.3 manifest 里的 data_file 结构

V2 的 data_file 结构(简化):

data_file: {
  content: int              // 0=DATA, 1=POSITION_DELETES, 2=EQUALITY_DELETES
  file_path: string
  file_format: string       // avro / orc / parquet
  partition: struct
  record_count: long
  file_size_in_bytes: long
  column_sizes: map<int, long>
  value_counts: map<int, long>
  null_value_counts: map<int, long>
  nan_value_counts: map<int, long>
  lower_bounds: map<int, binary>
  upper_bounds: map<int, binary>
  key_metadata: binary
  split_offsets: list<long>
  equality_ids: list<int>
  sort_order_id: int
  referenced_data_file: string   // 字段 ID 143,V2 可选
}

记住 contentreferenced_data_file,V3 的改动主要落在这里。


三、V3 核心特性一:Deletion Vectors(删除向量)

3.1 设计目标

V3 规范对 DV 提出了一个极强的不变量:

一个数据文件,在任意时刻,最多只有一个活跃的删除向量。

这一句话直接消灭了 V2 的"N 个 delete file"问题。读路径从"合并 N 个流"变成"读 1 个 bitmap"。

代价是什么?写的时候必须做 read-modify-write:新增删除时,要先读出旧 DV,或上新的位置,再写一个新的 DV。

这是一个经典的权衡:把复杂度从读侧移到写侧。对分析型负载(写少读多、或者写批量读高频)这是绝对划算的。

3.2 物理载体:Puffin 文件格式

DV 不存在 Parquet 里,而是存在 Puffin 文件里。

Puffin 是 Iceberg 自己定义的一个"辅助数据"容器格式,最早是为了存 Theta Sketch(近似去重计数)。它的结构非常简单:

+------------------------------------------+
|  Magic: "PFA1" (4 bytes)                 |
+------------------------------------------+
|  Blob 1 (任意字节)                        |
+------------------------------------------+
|  Blob 2                                  |
+------------------------------------------+
|  ...                                     |
+------------------------------------------+
|  Magic: "PFA1"                           |
|  FooterPayload (JSON, 可能被压缩)          |
|  FooterPayloadSize (4 bytes, LE)         |
|  Flags (4 bytes)                         |
|  Magic: "PFA1"                           |
+------------------------------------------+

Footer 里的 JSON 描述了每个 blob 的类型、偏移、长度、属性:

{
  "blobs": [
    {
      "type": "deletion-vector-v1",
      "fields": [],
      "snapshot-id": 8744736658442914487,
      "sequence-number": 15,
      "offset": 4,
      "length": 156,
      "properties": {
        "referenced-data-file": "s3://bucket/data/00000-1-abc.parquet",
        "cardinality": "1337"
      }
    }
  ]
}

关键设计:一个 Puffin 文件可以装多个 DV blob。这非常重要——如果一次 commit 删了 1000 个数据文件里的行,不需要写 1000 个小文件,可以打包成一个 Puffin。这直接缓解了对象存储的小文件问题。

3.3 DV blob 的字节布局

这是最硬核的部分。deletion-vector-v1 类型的 blob,内部结构是:

+-------------------------------------------------+
|  length: 4 bytes, big-endian                    |  <- magic + bitmap 的总长度
+-------------------------------------------------+
|  magic: 4 bytes, 0xD1 0xD3 0x39 0x64            |
+-------------------------------------------------+
|  bitmap: 64-bit Roaring Bitmap, portable 格式    |
+-------------------------------------------------+
|  CRC-32C: 4 bytes, big-endian                   |  <- 校验 magic + bitmap
+-------------------------------------------------+

几个值得注意的设计:

为什么有个冗余的 length? Puffin footer 里已经有 blob length 了。这里再存一份是为了让 DV blob 自描述——可以脱离 footer 单独解析,方便调试工具和跨系统兼容(Delta Lake 的 DV 也有类似自描述头)。

为什么是 CRC-32C 而不是 CRC-32? Castagnoli 多项式在现代 CPU 上有硬件指令(x86 SSE4.2 CRC32、ARM CRC32CB/CH/CW/CX),吞吐 10+ GB/s,查表法 CRC-32 只有 1GB/s 量级。读路径上要校验每个 DV,这个差距很实在。

为什么要求"不压缩"? 规范要求 DV blob 不设 compression-codec。Roaring Bitmap 本身就是压缩数据结构,再套一层 zstd 收益通常 <5%,却增加解压 CPU 和内存分配。

3.4 Roaring Bitmap:为什么它是这个场景的最优解

Roaring Bitmap 是这套设计的灵魂。理解它才能理解 DV 的性能特征。

朴素方案的问题:排序整数数组,删 100 万行就是 800 万字节,稀疏时还行、稠密时爆炸;朴素位图,一个 1 亿行的文件固定要 12.5MB,无论你删 1 行还是 1 亿行。

Roaring 的做法:分桶 + 自适应容器。

把 64 位整数空间切成高 48 位(桶键)和低 16 位(桶内偏移)。每个桶最多容纳 65536 个值,桶内根据基数自动选择三种容器之一:

容器类型适用基数存储方式空间
Array Container≤ 4096uint16[] 有序数组2 x n 字节
Bitmap Container> 409665536-bit 位图固定 8KB
Run Container连续段少(start, length)4 x runs 字节

4096 这个阈值是怎么来的? 数学上:array container 存 n 个值需要 2n 字节,bitmap container 固定 8192 字节。两者相等时 n = 4096。超过 4096 用 bitmap 更省,反之用 array 更省。这是一个精确的最优切换点,不是拍脑袋定的。

Run Container 的威力:如果你 DELETE FROM t WHERE dt < '2026-01-01',导致某个数据文件的前 500 万行全被删了,Run Container 只需要存一个 (0, 5000000) ——8 个字节搞定 500 万行的删除。

实际空间对比(一个 1000 万行的 Parquet 数据文件):

删除模式删除行数排序数组朴素位图Roaring
稀疏随机1,0008 KB1.25 MB~2 KB
中等随机100,000800 KB1.25 MB~200 KB
稠密随机5,000,00040 MB1.25 MB~1.25 MB
连续区间5,000,00040 MB1.25 MB~50 字节
全删10,000,00080 MB1.25 MB~50 字节

看最后两行——这就是为什么大批量删除在 V3 下几乎免费。

"portable" 格式是什么意思? Roaring 有两种序列化格式:Java 原生格式和 "portable"(可移植)格式。后者是跨语言标准(RoaringFormatSpec),Java / C++ / Go / Rust / Python 的实现都能读。Iceberg V3 强制要求 portable 格式,就是为了让 PyIceberg、iceberg-rust、Trino(Java)、DuckDB(C++)能互相读懂同一个 DV。这是"开放表格式"的题中之义。

3.5 manifest 层面的变化

V3 里,DV 在 manifest 中仍然登记为一个"delete file"条目,但字段语义有变:

data_file: {
  content: 1                        // 仍然是 POSITION_DELETES
  file_path: "s3://.../xxx.puffin"  // 指向 Puffin 文件
  file_format: "puffin"             // 新增的 format 值
  referenced_data_file: "s3://.../00000-1-abc.parquet"   // V3 中对 DV 为必填
  content_offset: 4                 // 字段 144:blob 在 puffin 中的起始偏移
  content_size_in_bytes: 156        // 字段 145:blob 长度
  record_count: 1337                // 被删除的行数(= bitmap 基数)
  ...
}

content_offset / content_size_in_bytes 是 V3 新增的字段。 它们的存在意味着:读引擎不需要读 Puffin footer,直接一次 ranged GET 就能拿到 DV blob。

这是个很实用的优化。想想 S3 的访问模型:

  • 无这两个字段:GET footer(要先读文件尾部找 footer size)-> 解析 JSON -> GET blob。至少 2-3 次往返
  • 有这两个字段:直接 GET Range: bytes=4-1591 次往返

对于一个扫描 1000 个数据文件的查询,这就是 2000 次多余的 S3 请求 vs 0 次的差别。

3.6 手写一个 DV 解析器

光看规范不过瘾,我们用 Python 从零解析一个 DV blob,验证上面所有细节。

"""
iceberg_dv_parser.py
从 Puffin 文件里解析 Iceberg V3 Deletion Vector
依赖:pip install pyroaring  (或者纯手写 roaring 解析,见下)
"""
import struct
import json
import zlib
from dataclasses import dataclass
from typing import List, Set

PUFFIN_MAGIC = b"PFA1"
DV_MAGIC = bytes([0xD1, 0xD3, 0x39, 0x64])


def crc32c(data: bytes) -> int:
    """CRC-32C (Castagnoli)。生产请用 google-crc32c(走硬件指令)"""
    poly, crc = 0x82F63B78, 0xFFFFFFFF   # 0x1EDC6F41 的反射
    for byte in data:
        crc ^= byte
        for _ in range(8):
            crc = (crc >> 1) ^ (poly if crc & 1 else 0)
    return crc ^ 0xFFFFFFFF


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


def read_puffin_footer(path: str) -> List[BlobMetadata]:
    """
    解析 Puffin footer
    尾部布局: Magic(4) | FooterPayload(N) | PayloadSize(4,LE) | Flags(4) | Magic(4)
    """
    with open(path, "rb") as f:
        assert f.read(4) == PUFFIN_MAGIC, "不是 Puffin 文件"
        f.seek(-16, 2)
        tail = f.read(16)
        payload_size = struct.unpack("<I", tail[0:4])[0]
        if tail[4] & 0x01:                # flags bit0 = footer 用 LZ4 压缩
            raise NotImplementedError("需要 lz4 库")
        f.seek(-(payload_size + 12), 2)
        meta = json.loads(f.read(payload_size).decode("utf-8"))

    return [BlobMetadata(b["type"], b["offset"], b["length"],
                         b.get("properties", {}))
            for b in meta["blobs"]]


def parse_dv_blob(raw: bytes) -> Set[int]:
    """
    解析 deletion-vector-v1 blob
    布局: length(4,BE) | magic(4) | roaring_bitmap | crc32c(4,BE)
    """
    assert len(raw) >= 12, "blob 太短"

    declared_len = struct.unpack(">I", raw[0:4])[0]
    magic = raw[4:8]
    assert magic == DV_MAGIC, f"magic 不匹配: {magic.hex()}"

    # declared_len 覆盖 magic + bitmap
    bitmap_end = 4 + declared_len
    bitmap_bytes = raw[8:bitmap_end]

    expected_crc = struct.unpack(">I", raw[bitmap_end:bitmap_end + 4])[0]
    actual_crc = crc32c(raw[4:bitmap_end])  # magic + bitmap
    assert expected_crc == actual_crc, (
        f"CRC 校验失败: expect={expected_crc:08x} actual={actual_crc:08x}"
    )

    return parse_roaring64_portable(bitmap_bytes)


def parse_roaring64_portable(data: bytes) -> Set[int]:
    """
    64-bit Roaring Bitmap,portable 序列化格式
    结构:
      uint64 (LE)  bucket_count
      for each bucket:
        uint32 (LE)  high 32 bits (bucket key)
        <32-bit roaring bitmap, portable format>
    """
    positions = set()
    pos = 0
    (bucket_count,) = struct.unpack_from("<Q", data, pos)
    pos += 8

    for _ in range(bucket_count):
        (high,) = struct.unpack_from("<I", data, pos)
        pos += 4
        low_vals, consumed = parse_roaring32_portable(data, pos)
        pos += consumed
        base = high << 32
        positions.update(base + v for v in low_vals)

    return positions


def parse_roaring32_portable(data: bytes, start: int):
    """
    32-bit Roaring portable 格式:
      cookie -> [run bitmap] -> 描述符数组(key, card-1) -> [offset 数组] -> 各容器数据
    """
    SERIAL_COOKIE = 12347          # 带 run container
    SERIAL_COOKIE_NO_RUN = 12346
    NO_OFFSET_THRESHOLD = 4

    pos = start
    (cookie,) = struct.unpack_from("<I", data, pos)
    pos += 4

    has_run = (cookie & 0xFFFF) == SERIAL_COOKIE
    if has_run:
        container_count = (cookie >> 16) + 1
        run_bitmap_size = (container_count + 7) // 8
        run_bitmap = data[pos:pos + run_bitmap_size]
        pos += run_bitmap_size
    else:
        assert cookie == SERIAL_COOKIE_NO_RUN, f"未知 cookie: {cookie}"
        (container_count,) = struct.unpack_from("<I", data, pos)
        pos += 4
        run_bitmap = b""

    # 描述符数组
    keys, cards = [], []
    for _ in range(container_count):
        k, c = struct.unpack_from("<HH", data, pos)
        pos += 4
        keys.append(k)
        cards.append(c + 1)   # 存的是 cardinality - 1

    # offset 数组(仅在容器数 >= 4 或无 run 时存在)
    if not has_run or container_count >= NO_OFFSET_THRESHOLD:
        pos += 4 * container_count

    values = []
    for i in range(container_count):
        key, card = keys[i], cards[i]
        base = key << 16
        is_run = has_run and bool(run_bitmap[i // 8] & (1 << (i % 8)))

        if is_run:
            (n_runs,) = struct.unpack_from("<H", data, pos)
            pos += 2
            for _ in range(n_runs):
                run_start, run_len = struct.unpack_from("<HH", data, pos)
                pos += 4
                # run_len 存的是 length - 1
                values.extend(base + run_start + j for j in range(run_len + 1))
        elif card <= 4096:
            # array container
            for _ in range(card):
                (v,) = struct.unpack_from("<H", data, pos)
                pos += 2
                values.append(base + v)
        else:
            # bitmap container: 固定 8192 字节 = 1024 个 uint64
            for w in range(1024):
                (word,) = struct.unpack_from("<Q", data, pos)
                pos += 8
                if word:
                    for b in range(64):
                        if word & (1 << b):
                            values.append(base + w * 64 + b)

    return values, pos - start


if __name__ == "__main__":
    import sys
    blobs = read_puffin_footer(sys.argv[1])
    with open(sys.argv[1], "rb") as f:
        for b in blobs:
            if b.type != "deletion-vector-v1":
                continue
            f.seek(b.offset)
            raw = f.read(b.length)
            deleted = parse_dv_blob(raw)
            ref = b.properties.get("referenced-data-file")
            card = b.properties.get("cardinality")
            print(f"数据文件: {ref}")
            print(f"  声明基数: {card}, 实际解析: {len(deleted)}")
            print(f"  前 20 个删除位置: {sorted(deleted)[:20]}")
            print(f"  blob 字节数: {b.length}, "
                  f"平均每删除行开销: {b.length / max(len(deleted),1):.3f} 字节")

跑一下真实文件的输出大概长这样:

数据文件: s3://lake/orders/data/00000-1-9f3a.parquet
  声明基数: 1337, 实际解析: 1337
  前 20 个删除位置: [42, 108, 209, 377, 512, 890, 1024, ...]
  blob 字节数: 2698, 平均每删除行开销: 2.018 字节

每行删除的元数据开销约 2 字节。对比 V2 的 position delete(Parquet 里每行是 path + long,即使字典编码后也要 8-12 字节,再加上 Parquet 自身的页头、footer、schema 开销,小文件的固定成本经常在几 KB),DV 的空间效率高一个数量级。

3.7 读路径的实际变化

V2 的读取伪代码:

// 简化版 V2 delete 应用逻辑
List<DeleteFile> deletes = findRelevantDeletes(dataFile);  // 可能几十上百个
PositionDeleteIndex index = new PositionDeleteIndex();
for (DeleteFile df : deletes) {
    try (CloseableIterable<Record> rows = openParquet(df)) {   // 网络 IO
        for (Record r : rows) {
            if (r.get("file_path").equals(dataFile.path())) {
                index.delete(r.get("pos"));
            }
        }
    }
}
// 然后扫描数据文件,对每行 index.isDeleted(rowPos)

V3 的读取伪代码:

// V3:最多一个 DV
DeleteFile dv = findDeletionVector(dataFile);   // 至多一个,或 null
PositionDeleteIndex index = (dv == null)
    ? PositionDeleteIndex.empty()
    : DeletionVectorUtil.loadRoaring(
          dv.path(), dv.contentOffset(), dv.contentSizeInBytes());  // 1 次 ranged GET

差别不只是"少读几个文件":可预测性(planning 阶段就能算出精确 IO 次数)、内存可控(上限就是 bitmap 大小,不会随历史累积失控)、可缓存(DV 不可变,(puffin_path, offset) 唯一确定内容,可放进 off-heap cache 跨查询复用;V2 的 delete file 组合是动态的,缓存命中率低)。

3.8 写路径的代价:read-modify-write

天下没有免费的午餐。V3 要求"一个数据文件最多一个 DV",意味着每次新增删除都要合并旧 DV。

// 概念性代码:向已有 DV 追加删除位置
RoaringBitmap64 existing = dv == null
    ? new RoaringBitmap64()
    : loadRoaring(dv);

existing.addAll(newDeletePositions);

// 写一个新的 Puffin blob(旧的在 expire 之后被 GC)
PuffinWriter writer = Puffin.write(outputFile).createdBy("spark-4.0").build();
writer.add(new Blob(
    "deletion-vector-v1",
    ImmutableList.of(),
    snapshotId,
    sequenceNumber,
    serializeDv(existing),
    null,                                   // 不压缩
    ImmutableMap.of(
        "referenced-data-file", dataFile.path().toString(),
        "cardinality", String.valueOf(existing.cardinality())
    )
));

两个实践后果:

  1. 高频小批量删除的成本上升。 每 10 秒删 5 行,每次都要读+写整个 DV;DV 长到 1MB 时,每个微批都要产生 1MB 写量。解法是攒批:把 Flink checkpoint interval 从 10s 调到 60s。
  2. 并发删除的冲突概率上升。 两个事务同时给同一数据文件加 DV,会在 commit 阶段冲突重试;V2 里两个 position delete file 可以共存不冲突。

所以 V3 的 DV 不是无脑更优,它优化的是"读多写少"或"批量写"的场景。对于极高频的行级删除,需要重新评估。

3.9 Equality Delete 怎么办

注意:V3 废弃的是 position delete file,不是 equality delete。后者仍然存在(Flink CDC 强依赖),但问题更严重——它没有 referenced_data_file,无法做文件级裁剪,读的时候必须对每个数据文件做 anti-join。

实践建议:在 Flink 写入侧就把 equality delete 转成 position delete(connector 的 upsert-enabled + 主键内部会做转换),或者定期 compaction 消化掉。


四、V3 核心特性二:Row Lineage(行血缘)

这是 V3 里我个人认为最被低估的特性。DV 是性能优化,Row Lineage 则是能力扩展——它让 Iceberg 从"表格式"往"可增量计算的数据底座"迈了一大步。

4.1 它解决什么问题

考虑一个场景:你有一张 Iceberg 事实表,上面挂了 10 个下游物化视图。上游更新了 1000 行,你希望下游只重算受影响的部分。

在 V2 里你怎么做?

  • SELECT * FROM t.changes 增量读?Iceberg 的 changelog scan 能给你 INSERT/DELETE,但它给不出"这一行是由哪一行更新来的"。一个 UPDATE 在存储层就是 DELETE + INSERT,两条记录之间没有身份关联。
  • 自己加一个业务主键列?可以,但主键可能变、可能有重复、可能为 NULL,而且每张表都要单独维护。

Row Lineage 在格式层给了每一行一个不可变的身份证。

4.2 两个隐藏列

V3 给每行加了两个保留字段(用极大的 field id,避免和用户列冲突):

字段名Field ID类型语义
_row_id2147483540long行的唯一标识,跨更新保持不变
_last_updated_sequence_number2147483539long该行最后被修改时的 sequence number

语义规则:

  • 新插入的行:_row_id 由系统分配新值,_last_updated_sequence_number = 当前 seq
  • 更新的行(DELETE+INSERT):新行继承旧行的 _row_id_last_updated_sequence_number 更新为当前 seq
  • 行被 compaction 重写但内容不变:两个字段都保持不变

最后一条极其重要。它意味着 _last_updated_sequence_number 是真正的"数据变更时间戳",不会被文件维护操作污染。你可以放心地写:

-- 拿到自 seq=100 以来真正变化的行
SELECT * FROM orders
WHERE _last_updated_sequence_number > 100;

而不用担心昨晚的 compaction 把整张表都"变更"了一遍。

4.3 精妙之处:不读数据就能分配 ID

朴素实现:写入时给每行分配一个自增 ID。问题是——这需要一个全局计数器,而且要把 ID 物理写进每个 Parquet 文件。前者是分布式系统的噩梦(每个 writer task 都要争抢),后者浪费空间(一列 int64,1 亿行就是 800MB)。

V3 的解法是分层预留 + 隐式推导

在三个层级引入了 row id 相关的元数据:

表元数据层(metadata.json):

{
  "format-version": 3,
  "next-row-id": 500000,        // 下一个可分配的 row id
  ...
}

快照层(manifest list 中的 snapshot 条目):

{
  "snapshot-id": 8744736658442914487,
  "sequence-number": 15,
  "first-row-id": 400000,       // 本快照分配的 row id 起始值
  "added-rows": 100000,         // 本快照新增行数
  ...
}

Manifest 层 / 数据文件层:

manifest_file: {
  first_row_id: 400000          // 该 manifest 覆盖的 row id 起始
  ...
}

data_file: {
  first_row_id: 400512          // 该数据文件第一行的 row id
  record_count: 1024
  ...
}

推导规则:如果一个数据文件的 _row_id没有被物理写入(这是常态),那么第 i 行的 row id = data_file.first_row_id + i

于是:

def compute_row_id(data_file, row_position, physical_row_id=None):
    """
    V3 row id 推导
    physical_row_id: 如果 Parquet 里真的存了 _row_id 列,用它
    """
    if physical_row_id is not None:
        return physical_row_id          # 显式值优先(用于 UPDATE 继承场景)
    if data_file.first_row_id is None:
        return None                     # V2 表或未启用血缘
    return data_file.first_row_id + row_position

零存储开销。绝大多数行(新插入、从未更新过的)不需要在 Parquet 里存 _row_id,一个 first_row_id + 行位置就够了。

只有在更新继承的场景下,新写的数据文件才需要真的物化 _row_id 列(因为要写入的是旧行的 ID,不是连续的)。Parquet 的字典编码 + RLE 对这种列压缩效果也很好。

4.4 分配算法:无锁的 ID 空间划分

commit 流程大致是:

1. Writer 侧统计本次要新增的行数 N(各 task 汇总)
2. 提交时读取当前 metadata 的 next-row-id = R
3. 本快照的 first-row-id = R,added-rows = N
4. 各 manifest / data file 按顺序切分 [R, R+N) 这个区间
5. CAS 更新 metadata: next-row-id = R + N
6. 如果 CAS 失败(并发提交),重新读 next-row-id,重算区间,重试

注意第 4 步:ID 区间的切分是在 commit 阶段做的,不需要重写数据文件(因为 ID 只存在于 manifest 元数据里)。这是整个设计能成立的关键。

冲突处理:两个并发提交都想要 [500000, 600000),只有一个能 CAS 成功。失败者重试时会拿到 [600000, 700000)。注意重试只需要改 manifest 的 first_row_id 字段,不用重写 Parquet——因为 Parquet 里根本没存这个值。

这里体现了一个很漂亮的分层设计原则:把易变的元数据留在最上层,让下层保持不可变

4.5 实战:用 Row Lineage 做增量物化视图

-- 建表时开启行血缘(V3 表默认开启)
CREATE TABLE prod.db.orders (
    order_id   BIGINT,
    user_id    BIGINT,
    amount     DECIMAL(18,2),
    status     STRING,
    updated_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(updated_at))
TBLPROPERTIES (
    'format-version' = '3',
    'write.delete.mode' = 'merge-on-read',
    'write.update.mode' = 'merge-on-read',
    'write.merge.mode'  = 'merge-on-read'
);

-- 查看行血缘元数据列
SELECT
    order_id,
    status,
    _row_id,
    _last_updated_sequence_number AS last_seq
FROM prod.db.orders
LIMIT 10;
order_id | status    | _row_id | last_seq
---------|-----------|---------|----------
  100001 | PAID      |       0 |        1
  100002 | PENDING   |       1 |        1
  100003 | SHIPPED   |       2 |        7   <- 被更新过
  100004 | PAID      |       3 |        1

现在做增量同步:

"""incremental_mv.py — 基于 row lineage 的增量物化视图刷新"""
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.appName("incremental-mv").getOrCreate()
CKPT   = "prod.meta.mv_checkpoints"
SOURCE = "prod.db.orders"
TARGET = "prod.mart.user_order_summary"
MV     = "user_order_summary"


def last_synced_seq() -> int:
    row = (spark.table(CKPT).filter(F.col("mv_name") == MV)
           .select("last_seq").first())
    return row["last_seq"] if row else 0


def current_seq() -> int:
    return (spark.table(f"{SOURCE}.snapshots")
            .orderBy(F.col("committed_at").desc())
            .first()["sequence_number"])


def refresh():
    last, cur = last_synced_seq(), current_seq()
    if cur <= last:
        print(f"无新变更 (last={last}, cur={cur})")
        return

    # 关键:只读真正变化的行。compaction 重写不会抬高 last_updated_seq,
    # 所以不会误触发全量刷新——这是 V2 时代做不到的
    changed = (spark.table(SOURCE)
               .filter(F.col("_last_updated_sequence_number") > last)
               .select("_row_id", "user_id", "amount", "status"))
    n = changed.count()
    print(f"检测到 {n} 行变更 (seq {last} -> {cur})")

    if n:
        changed.createOrReplaceTempView("changed_rows")
        # _row_id 就是天然的、跨更新稳定的主键
        spark.sql(f"""
            MERGE INTO {TARGET} t USING changed_rows s
            ON t.row_id = s._row_id
            WHEN MATCHED AND s.status = 'CANCELLED' THEN DELETE
            WHEN MATCHED THEN UPDATE SET
                t.user_id = s.user_id, t.amount = s.amount, t.status = s.status
            WHEN NOT MATCHED THEN INSERT (row_id, user_id, amount, status)
                VALUES (s._row_id, s.user_id, s.amount, s.status)
        """)

    spark.sql(f"""
        MERGE INTO {CKPT} t
        USING (SELECT '{MV}' AS mv_name, {cur} AS last_seq) s
        ON t.mv_name = s.mv_name
        WHEN MATCHED THEN UPDATE SET t.last_seq = s.last_seq
        WHEN NOT MATCHED THEN INSERT *
    """)


if __name__ == "__main__":
    refresh()

这段代码在 V2 时代是写不出来的。 你要么用业务主键(不可靠),要么全量重算(贵),要么用 changelog scan 自己拼 UPDATE 语义(复杂且易错)。

4.6 一个容易踩的坑

_row_id 只在单张表内唯一,不保证跨表、跨引擎语义,别把它当业务 ID 存到下游:表被 DROP 重建后 ID 会重置;migrate / add_files 导入的老数据可能没有 lineage;某些引擎的写入路径还不支持继承语义(UPDATE 会变成"新 ID")。升级前务必验证你用的引擎版本——这块实现成熟度参差不齐。


五、V3 核心特性三:Variant 类型

5.1 半结构化数据的世纪难题

每个做数据平台的人都遇到过这个问题:埋点日志里有一个 properties 字段,是 JSON,schema 随业务变化,有几百个可能的 key,但每条记录只用其中几个。

方案 A:存成 STRINGget_json_object(props, '$.utm_source'))。每次查询全量解析 JSON,CPU 爆炸;无法列裁剪(读一个字段得读整串);无法谓词下推(列统计对 JSON 字符串没意义);key 名每行重复存一遍。

方案 B:存成 MAP<STRING, STRING>。比 A 好一点(key 有字典编码),但丢失类型、丢失嵌套结构、依然无法列裁剪。

方案 C:全部展开成列。Schema 几百上千列且绝大多数是 NULL,Parquet footer 膨胀到几 MB,planning 变慢;新增 key 还要改 schema。

5.2 Variant:二进制自描述类型

V3 引入的 variant 类型,本质上是一个紧凑的二进制 JSON,由两部分组成:

variant := (metadata: binary, value: binary)

metadata 是一个字符串字典,存所有出现过的 key 名:

+----------------------------------------+
| header (1 byte)                        |
|   bits 0-3: version (= 1)              |
|   bit 4:    sorted_strings             |
|   bits 6-7: offset_size_minus_one      |
+----------------------------------------+
| dictionary_size (offset_size bytes)    |
+----------------------------------------+
| offset[0..dictionary_size]             |
+----------------------------------------+
| bytes: "user_id" "utm_source" "device" |
+----------------------------------------+

value 是实际的值,用一个 tag 字节区分类型:

header byte 的低 2 位 = basic type:
  0 = Primitive
  1 = Short String(长度 <= 63,直接内联)
  2 = Object
  3 = Array

Primitive 的子类型包括 null / boolean / int8/16/32/64 / double / decimal4/8/16 / date / timestamp / timestamp_ntz / float / binary / string 等。

核心优势:

  1. 类型保真:数字就是数字,不会退化成字符串。{"amount": 99.5} 存成 double,可以直接参与算术。
  2. Key 字典化:一个 Parquet row group 里的所有 variant 共享 metadata 字典的机会(取决于实现),key 名不重复存储。
  3. 无需解析即可跳过:object 的每个字段有偏移索引,取 $.utm_source 时可以二分查找 + 直接跳到偏移,不用扫描整个文档。
  4. Schema 自由:不需要预定义,新 key 直接写。

5.3 Variant Shredding:性能的杀手锏

光有紧凑编码还不够。真正的性能突破在 shredding(打散/物化)

思路是:把高频访问的字段,从 variant 里"抽"出来,物化成独立的 Parquet 列

在 Parquet 里,一个 shredded variant 列的物理布局大致是:

event_props (group)
+-- metadata      : BINARY              <- 字典
+-- value         : BINARY (optional)   <- 未被 shred 的剩余部分
+-- typed_value (group)                 <- 被 shred 出来的字段
    +-- utm_source (group)
    |   +-- value       : BINARY (optional)   <- 类型不匹配时的兜底
    |   +-- typed_value : BINARY (optional)   <- 强类型列
    +-- user_id (group)
    |   +-- value       : BINARY (optional)
    |   +-- typed_value : INT64 (optional)
    +-- ...

读取语义:对每个 shredded 字段,typed_valuevalue 最多只有一个非 NULL

  • 如果这一行的 utm_source 是字符串 -> 写进 typed_value
  • 如果这一行的 utm_source 意外是个数字(脏数据)-> 写进 value(保留原始 variant 编码)
  • 如果这一行没有这个 key -> 两个都是 NULL

这个设计非常聪明:它在"强类型列的性能"和"schema-less 的灵活性"之间做了完美的妥协。95% 的干净数据走强类型列(可以做谓词下推、向量化读取、字典编码),5% 的脏数据走兜底路径,不会导致整个 pipeline 失败。

性能收益:

查询JSON StringVariant(未 shred)Variant + Shredding
WHERE utm_source='wechat'全量解析所有行二分查 key + 解码 1 字段Parquet 列扫描 + 谓词下推
SELECT utm_source读整列 JSON读整列 variant只读 1 列
IO 量(1TB 表,取 1 个字段)~1TB~600GB~2GB

最后一行是关键:shredding 之后,取单个字段的 IO 从 TB 级降到 GB 级,因为 Parquet 的列存特性终于能发挥作用了。

5.4 实战代码

-- 建带 variant 列的表
CREATE TABLE prod.log.events (
    event_id    BIGINT,
    event_time  TIMESTAMP,
    event_name  STRING,
    props       VARIANT                -- V3 新类型
)
USING iceberg
PARTITIONED BY (hours(event_time))
TBLPROPERTIES (
    'format-version' = '3',
    -- 声明要 shred 的字段(具体属性名以引擎实现为准)
    'write.parquet.variant.shredding.enabled' = 'true'
);

-- 写入:直接给 JSON,引擎负责编码
INSERT INTO prod.log.events VALUES
  (1, TIMESTAMP '2026-08-01 10:00:00', 'page_view',
   PARSE_JSON('{"utm_source":"wechat","user_id":10086,"device":{"os":"iOS","ver":"18.2"}}')),
  (2, TIMESTAMP '2026-08-01 10:00:05', 'click',
   PARSE_JSON('{"utm_source":"douyin","user_id":10087,"element":"buy_btn"}'));

-- 查询:像访问结构体一样
SELECT
    event_id,
    props:utm_source::STRING       AS utm_source,
    props:user_id::BIGINT          AS user_id,
    props:device.os::STRING        AS os
FROM prod.log.events
WHERE props:utm_source::STRING = 'wechat'
  AND event_time >= TIMESTAMP '2026-08-01 00:00:00';

用 PyIceberg 检查 shredding 效果:

"""inspect_variant.py — 检查 variant 列的 shredding 布局与 IO 收益"""
import pyarrow.parquet as pq
from pyiceberg.catalog import load_catalog

catalog = load_catalog("prod", type="rest", uri="https://catalog.internal:8181")
tbl = catalog.load_table("log.events")
print(f"Format version: {tbl.metadata.format_version}")

# 拿一个数据文件,直接看 Parquet 物理列
sample = list(tbl.scan(limit=1).plan_files())[0].file.file_path
pf = pq.ParquetFile(sample)
print(f"采样文件: {sample}")
print("物理列:", *pf.schema.names, sep="\n  ")

# 按列统计压缩后大小,shredded 子列会单独出现
meta, col_sizes = pf.metadata, {}
for rg_idx in range(meta.num_row_groups):
    rg = meta.row_group(rg_idx)
    for c in range(rg.num_columns):
        col = rg.column(c)
        col_sizes[col.path_in_schema] = (
            col_sizes.get(col.path_in_schema, 0) + col.total_compressed_size)

total = sum(col_sizes.values())
print(f"\n总压缩大小: {total/1024/1024:.2f} MB\n列大小 Top 15:")
for path, size in sorted(col_sizes.items(), key=lambda x: -x[1])[:15]:
    mark = " <- shredded" if "typed_value" in path else ""
    print(f"  {size/1024:>10.1f} KB  ({size/total*100:5.2f}%)  {path}{mark}")

# 只读单个字段需要多少 IO
hit = sum(v for k, v in col_sizes.items() if "utm_source" in k)
print(f"\n只读 props.utm_source: {hit/1024:.1f} KB "
      f"({hit/total*100:.3f}% of table),节省 {(1-hit/total)*100:.2f}%")

典型输出:

总压缩大小: 128.44 MB
列大小分布(Top 15):
    52341.2 KB  (39.81%)  props.value
    18220.7 KB  (13.86%)  props.metadata
    11033.4 KB  ( 8.39%)  event_time
     8871.0 KB  ( 6.75%)  props.typed_value.user_id.typed_value <- shredded
     4102.9 KB  ( 3.12%)  props.typed_value.utm_source.typed_value <- shredded
     ...

只读 props.utm_source 需要 IO: 4110.3 KB
占全表比例: 3.126%
节省: 96.87%

5.5 Shredding 的选型建议

该 shred 什么:

  • 高频出现在 WHERE 里的字段(能吃到谓词下推)
  • 高频出现在 SELECT 里的字段(能吃到列裁剪)
  • 基数低、类型稳定的字段(字典编码效果好)

不该 shred 什么:

  • 出现率 <5% 的稀疏字段(一整列几乎全 NULL,白白增加 Parquet footer 的元数据开销)
  • 类型混乱的字段(一半是字符串一半是数字 -> 大量走 value 兜底路径,反而更慢)
  • 超大文本字段(shred 出来还是要读全量)

经验法则:shred 的字段数控制在 20-50 个。超过 100 个,Parquet footer 的膨胀会开始反噬 planning 性能。


六、V3 的其他改动(别忽略这些)

6.1 列默认值(Default Values)

V2 之前,ALTER TABLE ADD COLUMN 加的新列,对历史数据只能是 NULL。V3 引入了两个属性:

{
  "id": 5,
  "name": "region",
  "type": "string",
  "required": true,
  "initial-default": "unknown",   // 历史数据读出来的值
  "write-default": "cn-shanghai"  // 新写入时不指定则用这个
}

这是一个真正的刚需。 想加一个 NOT NULL 的列,V2 里你必须重写整张表。V3 里:

ALTER TABLE prod.db.orders
  ADD COLUMN region STRING NOT NULL DEFAULT 'unknown';
-- 元数据操作,秒级完成,不动一个数据文件

实现细节initial-default 存在 schema 里,读取时如果 Parquet 文件缺这一列,就用这个值填充。Iceberg 的列解析本来就是按 field ID 而非列名/位置,所以缺列是天然支持的场景,只是以前只能填 NULL。

6.2 纳秒时间戳

V3 新增 timestamp_nstimestamptz_ns

为什么需要?可观测性场景(trace span、性能 profile)里,微秒精度不够用。一个函数调用可能只有几百纳秒。

范围代价:int64 存纳秒,能表示的范围是约 ±292 年(1677 到 2262)。微秒精度是 ±29 万年。所以别拿 timestamp_ns 存生日或历史事件。

6.3 unknown 类型

一个"始终为 NULL"的占位类型。典型场景:从上游 schema 推断时遇到全 NULL 的列,先标成 unknown,后续再演进。因为它本来就没有值,可以安全 promote 成任何类型,让 schema 演进更灵活。

6.4 多参数转换

V2 的分区转换只接受单列(bucket(16, id)days(ts)),V3 允许多列:bucket(16, a, b)

这对 Join 优化很关键。 两张表都按 bucket(16, user_id, tenant_id) 分区,Join 时可以做 storage-partitioned join(分桶 join),完全跳过 shuffle。

6.5 Geometry / Geography

V3 加入地理空间类型,编码遵循 OGC 的 WKB 标准并携带 CRS 信息。geometry 是平面几何、支持自定义 CRS;geography 是球面几何、默认 WGS84、支持多种边缘插值算法(spherical / vincenty / thomas / andoyer / karney)。配合 Parquet 的 geospatial 统计(bounding box)可以做空间谓词下推,对 LBS、物流、地图业务是刚需。

另外 V3 在 manifest 层面强化了 key_metadata,为表级加密(Iceberg Encryption)铺路,但这块还在演进,生产使用需谨慎。


七、实战:从 V2 升级到 V3

7.1 升级操作本身很简单

-- 方式 1:ALTER TABLE
ALTER TABLE prod.db.orders SET TBLPROPERTIES ('format-version' = '3');

-- 方式 2:建表时指定
CREATE TABLE prod.db.new_orders (...) USING iceberg
TBLPROPERTIES ('format-version' = '3');

用 Java API:

Table table = catalog.loadTable(TableIdentifier.of("db", "orders"));
TableOperations ops = ((BaseTable) table).operations();
TableMetadata current = ops.current();
TableMetadata upgraded = current.upgradeToFormatVersion(3);
ops.commit(current, upgraded);

用 PyIceberg:

from pyiceberg.catalog import load_catalog

catalog = load_catalog("prod")
tbl = catalog.load_table("db.orders")

with tbl.transaction() as tx:
    tx.upgrade_table_version(format_version=3)

print(f"升级后: v{tbl.metadata.format_version}")

7.2 但是——升级是单向的

Iceberg 不支持降版本。 format-version 只能升不能降。

一旦升到 V3,写入侧会开始产生 Puffin DV、可能带 row lineage 元数据,任何不支持 V3 的读引擎会直接报错

所以升级前必须做兼容性审计。

7.3 兼容性检查清单

第一步:盘点所有读写这张表的系统。 不只是 Spark / Flink / Trino 这些你知道的,还有某个 BI 工具直连、数据科学团队的 Jupyter + PyIceberg,以及——最危险的——某个上古 Python 脚本用 pyarrow.dataset 直接读 Parquet 路径。

绕过 Iceberg 元数据直接读 Parquet 的代码,在 V3 下会读到已被 DV 删除的行,而且不报任何错。

第二步:检查各引擎的最低版本要求。

这块变化很快,务必查各自的官方文档确认当前状态。截至 2026 年年中的大致情况:

引擎V3 支持情况
Iceberg Java1.9+ 逐步完善,1.11 较成熟,以 release note 为准
Spark / Flink跟随 Iceberg 版本,需匹配 runtime jar;Flink 的 DV 写入路径需确认
Apache Hive4.2.0 起支持 DV / 列默认值 / Variant(官方明确宣布)
Trino较新版本支持读,写支持在推进,查 connector 文档
PyIceberg需较新版本,Variant 支持较晚
Snowflake2026 年已宣布 V3 支持
StarRocks / Doris / DuckDB视版本,查各自文档;DuckDB 扩展以只读为主

别信任何"我记得支持",动手写个探针测。

第三步:探针脚本。

"""v3_compat_probe.py — 升级前用一张小 V3 表验证各引擎兼容性
前提:沙箱表已插 100 行、删 10 行(从而产生 DV),期望读到 90 行
"""
import sys
from pyiceberg.catalog import load_catalog

PROBE = "sandbox.compat.v3_probe"
EXPECTED = 90


def check(name, fn):
    try:
        rows = fn()
    except Exception as e:
        print(f"✗ {name:16s} 报错: {type(e).__name__}: {e}")
        return False
    if rows == EXPECTED:
        print(f"✓ {name:16s} 读取正确 ({rows} 行)")
        return True
    print(f"✗ {name:16s} 行数错误 {rows},期望 {EXPECTED} ⚠ 可能未应用 DV!")
    return False


def probe_spark():
    from pyspark.sql import SparkSession
    return SparkSession.builder.getOrCreate().table(PROBE).count()


def probe_trino():
    import trino
    cur = trino.dbapi.connect(
        host="trino.internal", port=8080, user="probe").cursor()
    cur.execute(f"SELECT count(*) FROM iceberg.{PROBE.split('.', 1)[1]}")
    return cur.fetchone()[0]


def probe_pyiceberg(cat):
    return len(cat.load_table(PROBE).scan().to_arrow())


def probe_raw_parquet(cat):
    """模拟“绕过 Iceberg 直接读 Parquet”的老代码——必定读到已删行"""
    import pyarrow.dataset as ds
    paths = [f.file.file_path for f in cat.load_table(PROBE).scan().plan_files()]
    return ds.dataset(paths, format="parquet").count_rows()


if __name__ == "__main__":
    cat = load_catalog("prod")
    ok = [
        check("Spark", probe_spark),
        check("Trino", probe_trino),
        check("PyIceberg", lambda: probe_pyiceberg(cat)),
    ]
    print("\n--- 危险模式检测(下面失败是预期的)---")
    check("裸读 Parquet", lambda: probe_raw_parquet(cat))

    if not all(ok):
        print("\n❌ 存在不兼容引擎,不要升级生产表")
        sys.exit(1)
    print("\n✅ 所有引擎兼容,可以规划升级")

7.4 灰度升级策略

不要一次性升生产核心表。推荐四阶段:

  1. 沙箱验证(1-2 周):建 V3 沙箱表跑探针,覆盖所有下游读取路径。
  2. 影子表(2-4 周):选一张中等重要的表双写 V2 + V3,对比查询结果、P99 延迟、S3 请求数、compaction 成本。
  3. 边缘表(1-2 周):升级 3-5 张低风险表,观察一个完整业务周期(含月末批量、大促峰值)。
  4. 核心表:提前公告,升级后延长快照保留期。

关于回滚:格式不可降级,唯一的“回滚”是提前用 CREATE TABLE ... AS SELECT 做一份 V2 副本,出问题时切 catalog 指向副本。成本不低,所以前三个阶段的验证一定要做扎实。


八、性能优化:V3 时代的运维手册

8.1 Compaction 策略要重写

V2 时代的 compaction 逻辑是"delete file 太多了,重写数据文件把它们消化掉"。V3 下,DV 本身不是问题,问题变成了"DV 覆盖率太高"

新的判断标准:

"""v3_compaction_advisor.py — 基于 DV 覆盖率的 compaction 决策"""
import sys
from dataclasses import dataclass
from pyiceberg.catalog import load_catalog

TARGET_SIZE = 512 * 1024 * 1024


@dataclass
class FileHealth:
    path: str
    record_count: int
    deleted_count: int
    size_bytes: int

    @property
    def delete_ratio(self):
        return self.deleted_count / max(self.record_count, 1)

    @property
    def live_bytes(self):
        return int(self.size_bytes * (1 - self.delete_ratio))


def analyze(table_name):
    tbl = load_catalog("prod").load_table(table_name)
    files_df = tbl.inspect.files().to_pandas()
    deletes_df = tbl.inspect.delete_files().to_pandas()

    # V3 下每个 DV 都带 referenced_data_file,可以直接建立映射
    dv = {}
    for _, r in deletes_df.iterrows():
        ref = r.get("referenced_data_file")
        if ref:
            dv[ref] = dv.get(ref, 0) + int(r["record_count"])

    return [FileHealth(r["file_path"], int(r["record_count"]),
                       dv.get(r["file_path"], 0), int(r["file_size_in_bytes"]))
            for _, r in files_df.iterrows()]


def advise(table_name):
    files = analyze(table_name)
    n = len(files)
    phys = sum(f.size_bytes for f in files)
    live = sum(f.live_bytes for f in files)
    waste = phys - live

    high = [f for f in files if f.delete_ratio > 0.30]
    dead = [f for f in files if f.delete_ratio >= 0.999]
    small = [f for f in files if f.size_bytes < TARGET_SIZE * 0.25]

    print(f"=== {table_name} ===")
    print(f"数据文件总数: {n}")
    print(f"物理大小:     {phys/1024**3:.2f} GB")
    print(f"有效数据:     {live/1024**3:.2f} GB")
    print(f"浪费空间:     {waste/1024**3:.2f} GB ({waste/max(phys,1)*100:.1f}%)")
    print(f"高删除率文件(>30%): {len(high)}")
    print(f"全删文件:     {len(dead)}  <- 直接从 manifest 摘掉,零重写成本")
    print(f"小文件:       {len(small)}\n建议:")

    if dead:
        print(f"  1. 优先清理 {len(dead)} 个全删文件,释放 "
              f"{sum(f.size_bytes for f in dead)/1024**3:.2f} GB(几乎零成本)")
    if waste / max(phys, 1) > 0.20:
        print(f"  2. 浪费率 {waste/phys*100:.1f}% > 20%,对高删除率文件做 compaction")
    if len(small) > n * 0.3:
        print(f"  3. 小文件占比 {len(small)/n*100:.1f}%,做 bin-pack")
    if not (dead or small) and waste / max(phys, 1) <= 0.20:
        print("  表健康,暂不需要 compaction")
        return

    print(f"""
执行 SQL:
CALL prod.system.rewrite_data_files(
    table => '{table_name}', strategy => 'binpack',
    options => map(
        'delete-file-threshold', '1',
        'min-input-files', '5',
        'target-file-size-bytes', '{TARGET_SIZE}',
        'max-concurrent-file-group-rewrites', '10',
        'partial-progress.enabled', 'true',
        'partial-progress.max-commits', '10'));""")


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

关键洞察:全删文件是免费午餐。 一个 delete_ratio == 1.0 的数据文件,compaction 时不需要读它一个字节,直接从 manifest 里删掉引用即可。V3 的 DV 让"这个文件全删了"这个判断变得极其廉价(比较 dv.cardinality == data_file.record_count),而 V2 需要合并所有 delete file 才能确认。

在 CDC 场景下,全删文件的比例可能高达 20-30%(老数据被批量更新覆盖)。定期清理这部分是投入产出比最高的运维动作。

8.2 DV 的独立重写

V3 引入了一个新的运维需求:DV 自己也会碎片化

如果每次 commit 都产生一个新的 Puffin 文件,一段时间后你会有成千上万个小 Puffin。虽然每个只有几 KB,但对象存储的请求数还是上来了。

对策是定期把多个 DV 打包进一个 Puffin:

-- 重写 position delete / DV(具体过程名和参数以引擎版本为准)
CALL prod.system.rewrite_position_delete_files(
    table => 'db.orders',
    options => map(
        'rewrite-all', 'false',
        'target-file-size-bytes', '67108864',
        'min-file-size-bytes', '2097152'
    )
);

8.3 Manifest 优化

这一条 V2 V3 通用,但 V3 因为多了 DV 条目,manifest 会更多,更需要维护:

CALL prod.system.rewrite_manifests('db.orders');

判断标准:如果 planning 时间(从提交查询到开始读数据)超过 3 秒,大概率是 manifest 问题。

-- 检查 manifest 健康度
SELECT
    count(*)                        AS manifest_count,
    sum(added_data_files_count
        + existing_data_files_count) AS total_files,
    avg(added_data_files_count
        + existing_data_files_count) AS avg_files_per_manifest,
    sum(length)/1024/1024            AS total_mb
FROM prod.db.orders.manifests;

健康的 manifest 应该:

  • 每个 manifest 覆盖 500-5000 个文件(太少 = manifest 太碎,太多 = 裁剪不精细)
  • 单个 manifest 大小在 8-64MB

8.4 快照过期与孤儿文件

V3 下这两件事更重要,因为 Puffin 文件也会变成孤儿:

-- 过期老快照(保留 7 天)
CALL prod.system.expire_snapshots(
    table => 'db.orders',
    older_than => TIMESTAMP '2026-07-25 00:00:00',
    retain_last => 10,
    max_concurrent_deletes => 20
);

-- 清理孤儿文件(注意:一定要设 older_than,
-- 否则可能删掉正在写入的文件)
CALL prod.system.remove_orphan_files(
    table => 'db.orders',
    older_than => TIMESTAMP '2026-07-31 00:00:00',
    dry_run => true          -- 先 dry run!
);

血泪教训remove_orphan_files 不加 older_than 或设得太近,会删掉并发写入任务正在生成的临时文件,导致任务失败甚至丢数。默认 3 天,不要改小

8.5 表级参数调优

ALTER TABLE prod.db.orders SET TBLPROPERTIES (
    -- 目标文件大小
    'write.target-file-size-bytes' = '536870912',       -- 512MB

    -- 只对前 N 列存统计(列多时能显著缩小 manifest)
    'write.metadata.metrics.default' = 'counts',
    'write.metadata.metrics.column.order_id' = 'full',
    'write.metadata.metrics.column.updated_at' = 'full',

    -- manifest 自动合并
    'commit.manifest-merge.enabled' = 'true',
    'commit.manifest.target-size-bytes' = '16777216',   -- 16MB
    'commit.manifest.min-count-to-merge' = '100',

    -- 快照保留
    'history.expire.max-snapshot-age-ms' = '604800000', -- 7 天
    'history.expire.min-snapshots-to-keep' = '10',

    -- 删除模式
    'write.delete.mode' = 'merge-on-read',
    'write.update.mode' = 'merge-on-read',
    'write.merge.mode'  = 'merge-on-read'
);

write.metadata.metrics.default = 'counts' 这一条经常被忽视但收益巨大。 默认值是 truncate(16),即对每一列都存 min/max(截断到 16 字节)。如果你有 500 列,每个 data file 的 manifest 条目会有 500 x 4 个统计项。10 万个文件 = 2 亿个统计项,manifest 能膨胀到几 GB。

把默认改成 counts(只存 null_count / value_count),只对真正用于过滤的几列开 full,manifest 大小能降 80% 以上,planning 时间同步下降。


九、踩坑清单:那些文档不会告诉你的事

坑 1:升级后老引擎静默读到脏数据

前面提过,但值得再强调。绕过 Iceberg 元数据直接读 Parquet 路径的代码,在 V3 下会读到被 DV 删除的行,且不会报任何错。

排查方法:搜代码库里所有 pyarrow.datasetpd.read_parquetspark.read.parquet 且路径指向 Iceberg 表 warehouse 目录的用法。

坑 2:Row Lineage 在 add_files 导入时不连续

CALL system.add_files() 注册已有 Parquet 时,first_row_id 的分配可能和你预期不同,历史更新关系也无法恢复。Row Lineage 只对"从 V3 表开始就在里面"的数据有完整语义,存量迁入的数据把 lineage 当"从迁移那一刻开始"的新起点。

坑 3:Variant 的 shredding 是写时决定的

Shredding 配置改了只影响之后写入的文件,历史文件仍是老布局,想统一必须 compaction 重写。不同文件布局可能不同,读引擎按文件处理——这意味着 SELECT props:x 在老文件和新文件上性能可能差几个数量级。

坑 4:DV 的 read-modify-write 在高并发下的重试风暴

多个 writer 同时向同一批数据文件写删除会频繁 CAS 冲突。Iceberg 默认重试 4 次、指数退避,冲突严重时任务直接失败。调优:

ALTER TABLE prod.db.orders SET TBLPROPERTIES (
    'commit.retry.num-retries' = '10',
    'commit.retry.min-wait-ms' = '100',
    'commit.retry.max-wait-ms' = '60000',
    'commit.retry.total-timeout-ms' = '1800000'
);

但更好的做法是从架构上避免并发写同一张表:单 writer 或按分区分片。

坑 5:默认值不是“真的默认值”,纳秒时间戳支持最差

initial-default 只影响读取时缺列的填充。如果 Parquet 文件里这一列存在但值是 NULL,读出来就是 NULL,不会被替换成默认值。这个区别在数据质量排查时会让人抓狂。

另外 timestamp_ns 是 V3 里支持度最低的特性之一,很多引擎会静默降级成微秒或直接报错,用之前一定测。

坑 7:Puffin 生命周期依赖快照,元数据表 schema 变了

Puffin DV 的清理靠 expire_snapshots。如果快照保留 30 天,那 30 天内所有版本的 DV 都躺在存储上;对高频删除的表,Puffin 文件数可能比数据文件还多。用这条监控:

SELECT (SELECT count(*) FROM prod.db.orders.files)        AS data_files,
       (SELECT count(*) FROM prod.db.orders.delete_files) AS delete_files;

delete_files / data_files > 2 就该缩短快照保留期或增加 DV 重写频率。

另外 .files / .delete_files / .all_files 在 V3 下新增了 content_offsetcontent_size_in_bytesfirst_row_id 等列,任何 SELECT * 映射到固定 schema 的下游代码都可能挂掉。用显式列名。


十、横向对比:Iceberg V3 vs Delta Lake vs Hudi

2026 年这三家的格局已经比较清晰了,简单对比一下 V3 相关的能力:

能力Iceberg V3Delta LakeApache Hudi
删除向量Puffin + Roaring64Deletion Vector(更早,2023)有 log 文件机制,语义不同
行级血缘原生 _row_idRow Tracking(Databricks 侧较成熟)_hoodie_record_key 但是业务主键
半结构化类型Variant + ShreddingVariant(Spark 3.5+ 共享规范)依赖引擎
列默认值支持支持部分
引擎中立性最强(Trino/Flink/Spark/Snowflake/BigQuery…)Databricks 生态强,其他次之Spark/Flink 生态强
地理空间类型V3 原生无原生无原生
写入模型乐观并发 + 快照隔离乐观并发 + 事务日志时间线 + 多种表类型

一句话总结: Delta 在很多特性上先行一步(DV 早两年),但 Iceberg 的优势在于规范的开放性与引擎中立性——Variant 规范和 Delta/Spark 共享,Roaring portable 是跨语言标准,Puffin 是纯开放格式。Hudi 则一直更像“增量数据处理框架”而非纯表格式,流式 upsert 有独特优势,但生态广度不如前两者。


十一、总结:你现在该做什么

11.1 V3 到底值不值得升

值得升:CDC / 高频 upsert 的 MOR 表(DV 直接解决读放大);有埋点日志 / 半结构化数据(Variant + Shredding 是量级提升);需要增量物化视图或精确 CDC 下游(Row Lineage 是唯一优雅解);需要频繁 ADD COLUMN NOT NULL;地理空间业务。

暂时不必急:纯 append-only 的日志表(没删除,DV 用不上);Copy-on-Write 为主的表(DV 是 MOR 的优化);引擎版本被锁死的环境(风险大于收益)。

11.2 最后几句

V3 不是一次"加了几个新功能"的小版本。它做的是三件结构性的事:

  1. 删除向量——把 MOR 的读复杂度从 O(N 个 delete file) 降到 O(1),让 merge-on-read 真正变成生产可用的默认选项,而不是"用两周就得 compaction 救火"的权宜之计。

  2. 行血缘——第一次在格式层给了每一行稳定的身份。这打开的想象空间比 DV 大得多:增量物化视图、精确 CDC、数据血缘审计、行级权限追踪……这些以前需要在应用层用脆弱的 hack 实现的东西,现在有了扎实的地基。

  3. Variant——承认"schema 不是万能的"这个现实,用二进制编码 + shredding 在灵活性和性能之间找到了一个真正可用的平衡点。这比"要么全展开成列、要么存 JSON 字符串"的二选一进步了一代。

有意思的是,这三件事都不是"发明新东西"。DV 抄的是 Delta Lake,Roaring Bitmap 是十几年前的老技术,Variant 编码和 Spark 共享规范。Iceberg V3 的价值在于把这些东西以一种开放、跨引擎、可互操作的方式标准化了。

在一个 Snowflake、Databricks、BigQuery、Trino、StarRocks 各说各话的世界里,能让所有引擎读同一份 Roaring Bitmap、同一份 Variant 编码,这件事本身的价值可能比任何单个特性都大。

存储格式的战争在 2026 年基本结束了。接下来的战场是 catalog(REST Catalog 规范、Polaris、Unity Catalog、Gravitino 的博弈)和计算层的差异化。但至少在数据文件这一层,我们终于有了一个大家都认的标准。

对工程师来说,这是好事。少一层锁定,多一份选择权。


一手资料(规范变化很快,建议读原文):iceberg.apache.org/spec/(V3 章节)、iceberg.apache.org/puffin-spec/github.com/RoaringBitmap/RoaringFormatSpecgithub.com/apache/icebergformat/spec.md

本文所有字段名、属性名、过程名请以你实际使用的 Iceberg 版本文档为准。生产环境的每一个改动,都先在沙箱验证。

复制全文 生成海报 Iceberg 数据湖 湖仓一体 大数据 Spark 开源

推荐文章

deepcopy一个Go语言的深拷贝工具库
2024-11-18 18:17:40 +0800 CST
地图标注管理系统
2024-11-19 09:14:52 +0800 CST
JavaScript设计模式:装饰器模式
2024-11-19 06:05:51 +0800 CST
初学者的 Rust Web 开发指南
2024-11-18 10:51:35 +0800 CST
程序员茄子在线接单