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

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

Spark 与 pyiceberg 共用一个目录,重命名表,并找回已删除的表

在 TT Lab 中继续学习

目标

让 Spark 和 pyiceberg 共用 JDBC Catalog(一个 SQLite 文件),并通过提交前后的值,确认提交就是改掉 Catalog 中的一个指针。看到表的重命名和删除只发生在 Catalog 里、文件保持不变,并用一个 metadata 文件把删除的表恢复过来。

为什么重要

Iceberg 表的“真实状态”在 metadata 文件里,Catalog 只是把一个名称连到一个文件上的小表。然而,这张小表就是并发写入的裁判。两个写入方同时提交时,Catalog 用“只有我读到的指针仍未改变时才修改”这个条件,只接受其中一方。如果 Catalog 做不到这种原子替换,表就会损坏。 如果多个引擎看同一个 Catalog,Spark 创建的表可以由 Python 作业接着写入,而结果又可以由 Spark 再读取。另一方面,不知道名称、位置、文件是彼此独立的不同层,就会出事故。有人改了名称,就以为文件也被移动了,把旧路径删掉;有人以为 DROP TABLE 已经腾出了空间,就一直等着;也有人反过来,误删了表,就放弃了,认为永远找不回来。

步骤

  1. 用 /root/ice/cat/make.py(应用 ice-cat-make)创建 lake.cat.orders(format-version 2),并把 2026-03-01、03-02、03-03 三个文件一次性提交。
  2. 用 /root/ice/cat/list.py(pyiceberg,load_catalog("lake"))把命名空间和表列表写入 /root/ice/cat/out/tables.json。
  3. 用 /root/ice/cat/customers.py(pyiceberg)读取 /data/ice/customers.csv,创建 lake.cat.customers 并写入数据。
  4. 用 /root/ice/cat/join.py(应用 ice-cat-join)按 customer_id 连接两张表,创建各等级订单数的表 lake.cat.tier_counts(tier, orders)。
  5. 向 lake.cat.orders 提交 2026-03-04,同时把提交前的路径、提交后的路径和提交后的 previous 一栏写入 /root/ice/cat/out/pointer.json。
  6. 把 lake.cat.tier_counts 重命名为 lake.cat.tier_summary,并把重命名前后的 metadata 路径写入 /root/ice/cat/out/rename.json。
  7. 不带 PURGE 删除 lake.cat.tier_summary 之后,数一数剩下的数据文件个数,用删除前最后的 metadata 文件注册 lake.cat.tier_restored,并写入 /root/ice/cat/out/restore.json。
  8. 在 /root/ice/cat/report.md 中写出 ## 포인터、## 이름과 위치、## 지우기와 되살리기 三个小节。

参考

Spark 创建表——一次提交

创建 /root/ice/cat/make.py,应用名称为 ice-cat-make,创建 lake.cat 命名空间和 lake.cat.orders(六个列,'format-version' = '2'),并把 2026-03-01、03-02、03-03 三个文件用一次 append() 写入。

像 spark.read.csv([경로1, 경로2, 경로3], …) 这样传入列表,就会变成一个 DataFrame(占位符依次为三个文件路径)。提交只有一次,快照也就只有一个。评分器会检查第一个快照的行数是否是三个文件之和,以及摘要中是否有 Spark 留下的 engine-name。

用 pyiceberg 查看同一个 Catalog

用 pyiceberg 的 load_catalog("lake") 创建 /root/ice/cat/list.py,把命名空间和表列表按 {"namespaces": ["cat", …], "tables": ["cat.orders", …]} 的格式写入 /root/ice/cat/out/tables.json,并用 python3 list.py 运行。

pyiceberg 会从 ~/.pyiceberg.yaml 读取 lake Catalog 的配置(type: sql、uri: sqlite:////root/ice/catalog.db)。它是不启动服务器而打开同一个文件,所以 Spark 创建的表可以原样看到。名称是以元组形式给出的,请用点把它们连起来。

pyiceberg 创建表

用 /root/ice/cat/customers.py,以 pyarrow 读取 /data/ice/customers.csv(signup_date 为 date32),创建 lake.cat.customers 并写入数据。请使用 overwrite(),这样重新运行时行数也不会翻倍。

create_table_if_not_exists("cat.customers", schema=arrow_table.schema) 会把 pyarrow 的 schema 转换为 Iceberg 的 schema,并为每一列分配字段 ID。由 pyiceberg 写入的快照,摘要中没有 Spark 会留下的 engine-name——评分器就是用它来判断是谁写的。

Spark 连接两个引擎的表

用 /root/ice/cat/join.py,应用名称为 ice-cat-join,按 customer_id 连接 lake.cat.orders 和 lake.cat.customers,并把各等级(tier)的订单数做成 lake.cat.tier_counts(列为 tier、orders)(CREATE TABLE … AS SELECT)。

无论由谁写入,Iceberg 表都是同一份规范的 metadata 和 Parquet,所以引擎并不在意。评分器会用原始 CSV 直接计算各等级的订单数来对照(此时 orders 是 3 月 1–3 日的数据)。

一次提交,指针移动一格

记下 lake.cat.orders 当前的 metadata 路径,然后用 /root/ice/cat/append.py(日期参数)提交 2026-03-04,读取提交后的路径和 Catalog 的 previous_metadata_location,按 {"before", "after", "previous_after"} 写入 /root/ice/cat/out/pointer.json。

提交是先把新的 metadata 文件全部写好,然后在 Catalog 中用一次“只有当前值是 before 时才改成 after”的条件式更新来结束。所以提交之后,previous 一栏必须与提交之前的路径相同。可以用 jq -n --arg 把 shell 变量组合成 JSON。

名称只存在于 Catalog 中

执行 ALTER TABLE lake.cat.tier_counts RENAME TO cat.tier_summary(新名称不要带 Catalog lake),并把重命名之前(tier_counts)和之后(tier_summary)的 metadata 路径按 {"before", "after"} 写入 /root/ice/cat/out/rename.json。

JDBC Catalog 的重命名,是修改 iceberg_tables 中那一行的 table_name 的 UPDATE。表的位置(location)和文件保持不变,所以新名称的表仍然指向 …/cat/tier_counts/ 下的文件。评分器会检查两个路径是否相同,以及旧名称是否已从 Catalog 中消失。

用一个 metadata 文件恢复被删除的表

记下 lake.cat.tier_summary 的 metadata 路径,不带 PURGE 执行 DROP TABLE,数一数 /root/ice/warehouse/cat/tier_counts/data 中剩下的 Parquet 文件个数,再用 CALL lake.system.register_table(table => 'lake.cat.tier_restored', metadata_file => '<그 경로>') 把它恢复(占位符为该 metadata 路径)。把 {"metadata_file", "files_left"} 写入 /root/ice/cat/out/restore.json。

不带 PURGE 的 DROP 只删除 Catalog 中的一行。metadata 文件包含了 schema、快照和文件列表的全部内容,所以只要知道那个路径,就能把它重新注册到任何 Catalog 上(迁移 Catalog 或从备份恢复时,用的就是这种方法)。如果加了 PURGE,文件就会被删除,无法恢复。

Catalog 做什么、不做什么

在 /root/ice/cat/report.md 中写出 ## 포인터、## 이름과 위치、## 지우기와 되살리기 三个小节。第三节以数字写入第 7 步数出的剩余文件数。

请写出:在生产环境中更换 Catalog(例如 JDBC → REST),或者从备份恢复表时,必须迁移什么,什么可以保持原样。