Spark 与 pyiceberg 共用一个目录,重命名表,并找回已删除的表
目标
让 Spark 和 pyiceberg 共用 JDBC Catalog(一个 SQLite 文件),并通过提交前后的值,确认提交就是改掉 Catalog 中的一个指针。看到表的重命名和删除只发生在 Catalog 里、文件保持不变,并用一个 metadata 文件把删除的表恢复过来。
为什么重要
Iceberg 表的“真实状态”在 metadata 文件里,Catalog 只是把一个名称连到一个文件上的小表。然而,这张小表就是并发写入的裁判。两个写入方同时提交时,Catalog 用“只有我读到的指针仍未改变时才修改”这个条件,只接受其中一方。如果 Catalog 做不到这种原子替换,表就会损坏。 如果多个引擎看同一个 Catalog,Spark 创建的表可以由 Python 作业接着写入,而结果又可以由 Spark 再读取。另一方面,不知道名称、位置、文件是彼此独立的不同层,就会出事故。有人改了名称,就以为文件也被移动了,把旧路径删掉;有人以为 DROP TABLE 已经腾出了空间,就一直等着;也有人反过来,误删了表,就放弃了,认为永远找不回来。
步骤
- 用 /root/ice/cat/make.py(应用
ice-cat-make)创建lake.cat.orders(format-version 2),并把 2026-03-01、03-02、03-03 三个文件一次性提交。 - 用 /root/ice/cat/list.py(pyiceberg,
load_catalog("lake"))把命名空间和表列表写入 /root/ice/cat/out/tables.json。 - 用 /root/ice/cat/customers.py(pyiceberg)读取
/data/ice/customers.csv,创建lake.cat.customers并写入数据。 - 用 /root/ice/cat/join.py(应用
ice-cat-join)按customer_id连接两张表,创建各等级订单数的表lake.cat.tier_counts(tier, orders)。 - 向
lake.cat.orders提交 2026-03-04,同时把提交前的路径、提交后的路径和提交后的 previous 一栏写入 /root/ice/cat/out/pointer.json。 - 把
lake.cat.tier_counts重命名为lake.cat.tier_summary,并把重命名前后的 metadata 路径写入 /root/ice/cat/out/rename.json。 - 不带 PURGE 删除
lake.cat.tier_summary之后,数一数剩下的数据文件个数,用删除前最后的 metadata 文件注册lake.cat.tier_restored,并写入 /root/ice/cat/out/restore.json。 - 在 /root/ice/cat/report.md 中写出
## 포인터、## 이름과 위치、## 지우기와 되살리기三个小节。
参考
- Catalog 可以直接用
sqlite3 /root/ice/catalog.db "select * from iceberg_tables"查看。ice-loc <ns>.<표>会打印当前的 metadata 路径(占位符依次为命名空间与表名)。 - pyiceberg 的配置在
~/.pyiceberg.yaml中(Catalog 名称lake,type: sql)。Spark 一侧的配置在/opt/spark/conf/spark-defaults.conf。 - pyiceberg 第一次打开 Catalog 时会打印“v0 schema”警告。意思是它是 Java 的 JdbcCatalog 创建的旧形态(没有视图一栏),对写入表没有影响。
RENAME TO之后的新名称不要带 Catalog(cat.tier_summary)。如果写成lake.cat.tier_summary,Spark 会去找名为lake.cat的命名空间,然后以 NoSuchNamespaceException 中止(实测)。- 常见错误:在第 7 步使用
DROP TABLE … PURGE——连文件都被删掉,无法恢复。如果弄成了这样,请从第 4 步重新开始。 - 官方文档:JDBC Catalog · Spark DDL · Spark Procedures — register_table · pyiceberg — SQL Catalog
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),或者从备份恢复表时,必须迁移什么,什么可以保持原样。