TT Lab
开始
学习 学习路径 课程

湖仓表格式 — 从元数据理解 Apache Iceberg

一张表,多个引擎 — 规范是契约,目录是会合点

在 TT Lab 中继续学习

一句话总结

Iceberg 表是按规范写入的 metadata 和数据文件,所以 Python 可以往 Spark 创建的表里写入,DuckDB 可以读取。不过,这个约定只有在所有引擎都通过同一个 Catalog 找到“当前”metadata,并理解同样的规范功能和类型时,才能守住。

为什么需要多个引擎

批处理转换用 Spark 更方便,小规模加载和检查脚本用 Python 更方便,分析师的临时查询用 DuckDB 更方便。过去每个引擎各自放一张表,再互相复制。副本总会出现偏差,最后就会争论哪一边才是真的。Multi-Engine Support 文档把 Iceberg 介绍为任何处理引擎都能使用的开放标准。这是因为表不是由文件格式,而是由规范所规定的 metadata 来定义的。

同一份文档也展示了条件。Spark 和 Flink 的运行时 jar 是按引擎版本分别发布的,不在支持列表中的版本就无法写入。本课程的实验镜像之所以与仓库中的其他 Spark 课程(4.2)不同,使用 Spark 4.1.3,原因就在这里——Iceberg 1.11.0 的运行时只到 4.1。

工作原理——碰头之处与偏差之处

Catalog 是各引擎碰头的地方。在实验环境里,Spark 和 pyiceberg 打开的是同一个 SQLite Catalog。两者都通过条件式替换来提交,所以不会覆盖对方的提交。DuckDB 的 iceberg 扩展提供两种方式。指向 metadata 直接读取的方式不需要 Catalog,而且是只读的,如果还想写入,就要 ATTACH 一个 REST Catalog。

直接读取会固定在那一刻。metadata 文件一旦写入就不会改变,所以 iceberg_scan('…/00003-….metadata.json') 无论什么时候执行,返回的都是当时的表。即使有人在那之后又提交了,也不会有任何警告。DuckDB 文档指出,根据文件名猜测“最新”版本的功能可能破坏 ACID,所以默认是关闭的。要读取当前的内容,每次都必须重新从 Catalog 找到当前路径。

类型出现偏差的地方。规范的原始类型区分不带时区的 timestamp 和以 UTC 存储的 timestamptz。Spark 的 TIMESTAMP 会变成 timestamptz。所以,如果在 Python 中想以不带时区的时间写入同一个列,pyiceberg 会以 schema 不匹配为由拒绝(在实验镜像中已确认)。给时间加上 UTC 才是对的。还要记住,日期边界也是按 UTC 切分的。

功能出现偏差的地方。格式版本会在旧的读取一方无法正确读取新功能时提升。规范指出,为了避开引擎尚未实现的功能,可以继续以旧版本写入。在使用版本 3 的 deletion vector 或新类型之前,必须确认读取这张表的所有引擎都支持。不认识删除文件的引擎,会把已删除的行返回来。

配合良好的地方。列是按字段 ID 选取的,所以在一个引擎中重命名列之后,其他引擎也能用新名称读取旧文件的值。分区值在清单中,所以不管是谁写的,都会按同样的条件跳过文件。

是谁写的——快照摘要

引擎一多,总会需要知道“这次提交是从哪里来的”。快照摘要就是线索。在实验镜像中可以看到,Spark 会在摘要中留下 engine-name(spark)、engine-version 和 app-id,而 pyiceberg 0.12 不会留下。摘要是由写入的一方填写的,所以各引擎并不相同,这一点也要一并记住。

在现场相遇的样子

仪表板停在了昨天的数字。有人把 metadata.json 的路径写死在了 BI 工具里。表每天都在提交,那个工具读取的却还是第一天的表。应让工具通过 Catalog 读取,或者改成每次都去找当前路径。

Python 加载因 schema 错误而中止。从 CSV 读来的时间是不带时区的值。加上时区就行了。对这个错误要心存感激——如果悄悄写进去,就会混入偏差九个小时的值。

接入新引擎。要先确认:它是否支持 Catalog,能读取哪种格式版本和删除格式,怎样处理 timestamptz。三者只要有一个对不上,结果就会悄无声息地出错。

实际工作中真正重要的事

下一项实验要做什么

让 Spark 创建按日期划分的表并写入一周的数据,再由 pyiceberg 以带时区的时间多写一天。用 DuckDB 统计各地区的合计,用 pyiceberg 做条件读取,然后记下旧的 metadata 路径,让 Spark 再提交一天,分别用 DuckDB 读取旧路径和新路径。在 Spark 中重命名一个列,确认另外两个引擎是否用新名称读取,最后让 Spark 汇总三个引擎写入的全部行,生成按日汇总表。