Apache Iceberg Variant 类型实战指南:在 Iceberg v3 中存储任意形状的半结构化数据

原创2026-09-24 23:32:13294 阅读
文章标签:数据湖大数据数据存储

Apache Iceberg Variant 类型实战指南:在 Iceberg v3 中存储任意形状的半结构化数据

本文围绕 Apache Iceberg v3 新增的 variant(Variant)类型展开,系统讲解它要解决的半结构化数据存储难题、在 Iceberg 类型体系中的定位、底层的 Parquet/Avro/ORC 物理编码、统计与数据跳过机制,以及跨 Spark/Flink/Arrow 引擎的读写实操。读完本文,你将掌握如何在 Iceberg v3 表中用一个 variant 列承载结构不断变化的 JSON 类数据,并理解其与固定 schema 的取舍边界。

Variant 是 Iceberg v3 引入的一种"半结构化类型"(semi-structured type):单个列可以存放形状任意、持续演化的值,以紧凑的二进制形式存储,并被各引擎一致地读写。Parquet 项目已经定义了 Variant 类型及其二进制编码,因此 Iceberg 的增量工作集中在:Variant 如何融入表的 schema、数据文件、快照与统计信息体系。本文即以此为主线展开。

Variant 要解决的问题:固定 schema 表格式与半结构化数据的矛盾

先看一组典型的、形状随时间变化的事件数据:

{"event": "view", "page": "/pricing", "session": "s-8f21"}
{"event": "signup", "session": "s-8f21", "plan": "pro", "price": 29.00}
{"event": "view", "page": "/docs", "session": "s-3c07", "referrer": {"source": "search", "term": "iceberg variant"}}

三条记录字段各不相同,且后续还可能冒出新的字段。在没有 Variant 的情况下,Iceberg 表只能用两种方式存储这类数据,且各有明显缺陷:

  • 把 JSON 存成字符串。足够灵活,但读取单个字段必须先解析整段文本;同时 JSON 的类型系统很单薄:时间戳只是字符串,数值的精度也不明确。
  • 采用僵化的扁平 schema。查询很快,但每新增一个字段就是一次 schema 迁移;稀疏的、一次性出现的字段还会浪费大量存储空间。

Variant 像 JSON 一样灵活,却以紧凑、带类型的二进制来编码数据。值保留其原生类型:时间戳仍是时间戳,decimal 仍是精确的 decimal,而不会塌缩成 JSON 的字符串和普通数字。Iceberg schema 中该列只声明为 variant,不声明内部的字段和类型,因此每一行的结构可以互不相同,新增一个字段也无需任何表 schema 迁移。

Variant 在 Iceberg 类型系统中的定位

Variant 是在 v3 规范中加入 Iceberg 类型系统的。规范将其单独归为一类:variant 既不是原始类型(primitive type),也不是嵌套类型(nested type)。它比原始类型更丰富,却又不像 struct 或 list 那样有固定的、被声明的形状。这一点在 Iceberg 表规范 的"半结构化类型"一节有明确表述:variant 的结构和数据类型在表内或文件内的各行之间不必保持一致。

一个 Variant 值类似 JSON,但原始类型集合更宽,包含 date、timestamp、timestamptz、binary、decimal 等。它还可以嵌套:

  1. Variant array(Variant 数组):Variant 值的有序集合。与 Iceberg 的 list 不同,其元素不限于单一元素类型。
  2. Variant object(Variant 对象):以字符串为键、值为 Variant 值的字段集合。与 Iceberg schema 中的 struct 列不同,其字段不是一组固定命名的类型化列。

在 Iceberg schema 中,Variant 列以类型名 variant 声明,在表元数据中存储为普通字符串 "variant"(而非 struct 或 list 那样的嵌套对象)。在 Java API 侧,它对应 Types.VariantType:一个单例类型,typeId() 返回 TypeID.VARIANT,toString() 返回 "variant",并实现了 isVariantType() / asVariantType() 的访问方法——这正是规范中"声明为普通字符串"的落地实现。

由于形状不固定,Variant 列受一系列约束:不能与其他类型互相提升(promotion),不能用于分区,不能作为标识字段(identifier field),且默认值必须为 null。规范 format/spec.md 明确:"unknown、variant、geometry、geography 类型的所有列必须默认为 null,非 null 的 initial-default 或 write-default 均无效。"

Variant 如何存储:metadata 与 value 的二进制对

Iceberg 不自行定义 Variant 的二进制编码,类型与编码直接来自 Apache Parquet 项目的 Variant 编码规范(当前支持 V1),Iceberg 原样使用。因此一个 variant 列在 Parquet 中映射为一个含两个二进制字段的 group:

optional group payload (VARIANT(1)) {
  required binary metadata;
  required binary value;
}
  • metadata:存放该值用到的字段名词典,使字段名不必随数据内联重复出现。
  • value:存放编码后的数据本体——一个标量、数组或对象。数组和对象会为每个元素记录 field_offset(该元素值的起始字节偏移),对象还会为每个字段记录 field_id(指向 metadata 词典的下标)。

Variant 列本身和其他列一样通过 Iceberg 的 field ID 寻址,但它的 metadata 与 value 子字段是按名称访问的——这一点对后面的 shredding(拆分)至关重要。规范中的文件格式映射表(format/spec.md、format/spec.md、format/spec.md)给出了完整对应关系:

文件格式 variant 列的物理表示 备注
Parquet 含 metadata、value 两字段的 group,注解 VARIANT 支持 shredding
Avro 含 metadata、value 两字段的 record 不支持 shredding
ORC 含 metadata、value 两字段的 struct,iceberg.struct-type=VARIANT 不支持 shredding

也就是说,同一个 Variant 会映射进 Iceberg 支持的每一种文件格式:Parquet group、Avro record、ORC struct,内部都承载 metadata + value 这一对。在 Avro 和 ORC 中,Variant 始终是单个未拆分的 metadata + value 对。另外,Variant 的单值(single-value)二进制序列化在规范中是不被支持的(format/spec.md),这与它的多布局特性直接相关。

一个列,多种物理布局

Variant 的结构在行与行之间、文件与文件之间并不一致,但该列的 Iceberg 类型始终是 variant,无论有多少种形状流过它。数据内部新增或删除字段不需要 Iceberg schema 变更,只改变 Variant 值本身,以及(可能)每个文件各自的物理布局。

同一个逻辑列在每个数据文件中可以有不同的布局方式。在 Parquet 中,一个文件可能以未拆分的 metadata + value 对存储,而另一个文件可能把热点字段 shred 进专门的类型化列。两个文件携带相同的 Variant field ID,读取方会自行协调遇到的不同布局:

payload  (one Variant column, one field ID)
├─ data file A, unshredded:  metadata + value
└─ data file B, shredded:    metadata + value + typed_value.event, typed_value.country

快照不会改变这一点。每个快照记录的是写入时生效的 schema;由于 Variant 列在 schema 演化中保持其 field ID 不变,时间旅行(time travel)可以把所有文件都经由同一个列读回来。

统计与数据跳过

Variant 是 Iceberg schema 中的一列,因此表的 manifest 可以为其携带统计信息——这正是规划(planning)阶段跳过文件的依据。Iceberg 为 Variant 列记录 value 计数与 null 计数;按字段的 lower/upper bounds 则是可选的,存储为一个 Variant 对象,其键是各字段的归一化 JSON 路径。

关于 bounds 的细节,规范 format/spec.md 有更精确的说明:

  • bounds 对象的值是各字段的原始 Variant 下界/上界;是否包含某个字段的 bounds 是可选行为,且 upper/lower bounds 必须具有相同的 Variant 类型。
  • 字段路径采用归一化 JSON path 格式,例如 $['location']['latitude'] 或 $['user.name'];特殊路径 $ 表示 Variant 根节点本身的 bounds(适用于整个值都是同质原始类型、比如纯字符串的情况)。
  • 只有当字段的全部非 null 值类型一致(允许混入 null)时才允许写 bounds;例如某个 measurement 字段若同时混有 int64 与字符串值,就不能为其写 bounds。
  • 序列化方式是:编码后的 v1 metadata 与编码后的 bounds 对象拼接(format/spec.md)。

这里存在一个工程取舍:对于未拆分的 value,计算某字段的 bounds 意味着要读取原始值字节并解析,因此 writer 通常只记录计数。当字段被 shred 成独立的类型化列后,其 Parquet 统计信息可以直接获得,Iceberg 便能记录它的 bounds。在 MetricsConfig.java 的默认 metrics 收集实现中,variant 类型的访问器返回 null,即默认不把 Variant 列纳入按列收集的 metrics 集合,这与"未拆分时 bounds 计算成本高、writer 通常只记计数"的取舍方向一致。

一旦 manifest 中有了这些 bounds,针对相应字段的谓词就可以在读取任何数据之前剪除整个文件:如果文件记录的该字段范围不可能命中,Iceberg 会在规划阶段直接跳过它;而没有记录 bounds 的字段,则交由引擎读取后再过滤。

跨引擎读写

  • Apache Spark 4.0 与 4.1:可创建含 VARIANT 列的表并读写。读取支持在 Iceberg 1.10.0 中落地于 Spark 4.0;Iceberg 1.11.0 增加了 Spark 4.1 支持,以及在 Spark 4.0 与 4.1 上的 shredded Variant 写入。
  • Apache Flink 2.1:Variant 支持在 Iceberg 1.11.0 中加入,目前仅覆盖未拆分(unshredded)的 Variant;shredded 写入支持已合并,预计在后续版本中发布。

由于这些引擎使用同一个 Iceberg Variant 类型,由一个引擎写入的值在其它引擎中读回是完全一致的。Apache Arrow 在内存中通过其规范扩展类型 arrow.parquet.variant 承载同一个 Variant 值,让引擎之间无需特殊处理即可交换。

使用 Variant:Spark SQL 实操

用 Spark SQL,你可以把异构事件存进单个 VARIANT 列,再用 variant_get 读回字段。注意 Variant 是 v3 类型,表必须使用 format version 3:

-- Variant 是 v3 类型,因此表必须声明为 format version 3
CREATE TABLE events (id BIGINT, payload VARIANT)
USING iceberg
TBLPROPERTIES ('format-version' = '3');

-- 把不同形状的事件插入同一列
INSERT INTO events VALUES
    (1, parse_json('{"event": "login", "country": "US"}')),
    (2, parse_json('{"event": "purchase", "country": "UK", "amount": 99}')),
    (3, parse_json('{"event": "login"}'));

-- 用 variant_get(column, path, type) 从 Variant 中读出字段
SELECT
    id,
    variant_get(payload, '$.event', 'string')   AS event,
    variant_get(payload, '$.country', 'string') AS country,
    variant_get(payload, '$.amount', 'int')     AS amount
FROM events;

查询为每个事件返回一行,字段缺失处为 null:

id event country amount
1 login US null
2 purchase UK 99
3 login null null

第 3 行只有 event,因此 country 和 amount 读回为 null。Variant 不要求每行共享同一结构,异构事件得以共存于一个列中。

何时使用 Variant

当你不掌握数据的形状,或者它的变化速度超过你愿意演进 schema 的速度时,Variant 是正确的选择:

  • 事件与点击流数据:每种事件类型携带不同的字段集合,新字段随时间不断出现。
  • 应用与服务日志:结构化但异构的 payload。
  • IoT 与传感器遥测:每种设备型号上报各自的读数。
  • 第三方 API 响应与 webhook:schema 由别人拥有,可随时变更而无需通知。
  • 稀疏属性:否则会变成一张充满 null 列的宽表。

Variant 并不是已知 schema 的替代品。当某个字段出现在每一行上、且你经常查询它时,常规的类型化列(或者固定嵌套形状的 struct)更简单也更高效。一种常见模式是:把稳定、高频查询的字段保留为普通列,把变化的部分放进一个 Variant 列。

下一步:shredding

到目前为止,一个 Variant 列就是单一的 metadata + value 对。读取某个字段仍意味着解码 metadata 词典、再从那单个 value blob 中取出该字段——因此查询一个字段要读取整个 value 列。这正是 shredding 优化的用武之地:writer 选择把哪些字段抽取为各自的类型化列、并为每个字段推断类型,不匹配的值回退到未类型化的 value,reader 再从两者重构出原始 Variant。Iceberg 仓库中已能看到相关配置入口:TableProperties.java 定义了 write.parquet.shred-variants(默认 false)与 write.parquet.variant-inference-buffer-size(默认 100)。shredding 仅适用于 Parquet,其具体工作原理将在后续文章中展开。

进一步阅读

Iceberg 的 Variant 支持仍在持续演进,社区也通过官方邮件列表与 Slack 接受贡献与反馈。由于 Variant 涉及 v3 规范,使用前请确认你的引擎(Spark 4.0/4.1、Flink 2.1 及以上)与 Iceberg 版本(1.10.0 起支持 Spark 4.0 读取、1.11.0 起支持 Spark 4.1 与 Flink 2.1)满足相应要求。

登录后查看全文
iceberg