Spark 建表、Python 写入、DuckDB 读取 — 三个引擎共用一张表
目标
让 pyiceberg 往 Spark 创建的表里再写入一天的数据,让 DuckDB 和 pyiceberg 读取这张表,再由 Spark 汇总所有引擎写入的行。在此过程中,确认抱着旧 metadata 路径读取的“陈旧指针”陷阱,以及一个引擎的 schema 变更在其他引擎中是什么样子。
为什么重要
表格式的价值在于不必挑选引擎。批处理用 Spark,小规模加载用 Python,临时分析用 DuckDB,它们看到的都是同样的 metadata 和同样的 Parquet 文件。然而这个约定是有条件的。所有引擎都必须通过同一个 Catalog 找到“当前”metadata,并且必须支持同样的规范功能(格式版本、删除文件、类型)。 直接传入 metadata.json 路径来读取的工具虽然方便,但会固定在那个路径所指向的那一刻。如果有人在那之后提交了,你读到的还是旧表,却没有任何警告。类型也要当心——Spark 的 TIMESTAMP 就是 Iceberg 的 timestamptz,所以传入不带时区时间的 Python 代码,会被 schema 检查拦下。反过来,像重命名列这样靠字段 ID 解决的变更,在所有引擎中都原样可见。
步骤
- 用 /root/ice/eng/spark.py(应用
ice-eng-spark)按days(order_ts)划分创建lake.eng.orders,并把 2026-03-01 到 03-07 的七个文件一次性提交。 - 用 /root/ice/eng/py_append.py(pyiceberg)把 2026-03-08 追加到同一张表。
- 用 /root/ice/eng/duck.py(DuckDB)读取当前 metadata,把各地区的
amount合计写入 /root/ice/eng/out/duck_region.json。 - 用 /root/ice/eng/py_read.py(pyiceberg)把
region = 'seoul'且order_ts >= 2026-03-05的行数写入 /root/ice/eng/out/py_count.json。 - 记下当前的 metadata 路径,用 /root/ice/eng/append.py(应用
ice-eng-append)提交 2026-03-09,再用 /root/ice/eng/stale.py 分别以 DuckDB 读取旧路径和新路径,写入 /root/ice/eng/out/stale.json。 - 用 /root/ice/eng/rename.py(应用
ice-eng-rename)把region改为area,并用 /root/ice/eng/columns.py 把 pyiceberg 和 DuckDB 看到的列名写入 /root/ice/eng/out/rename.json。 - 用 /root/ice/eng/daily.py(应用
ice-eng-daily)创建每天订单数和金额合计的表lake.eng.daily(d, orders, amount)。 - 在 /root/ice/eng/report.md 中写出
## 한 표, 세 엔진、## 낡은 포인터、## 이름 바꾸기三个小节。
参考
- DuckDB 在
con.execute("LOAD iceberg")之后,用iceberg_scan('<metadata.json 경로>')读取(占位符为 metadata.json 路径)。路径由ice-loc eng.orders打印出来。扩展已提前放进了镜像(没有互联网)。 - 快照是谁写的,可以从摘要中看出来——Spark 会留下
engine-name: spark和app-id,而 pyiceberg 不会留下(SELECT summary FROM lake.eng.orders.snapshots)。 - 常见错误:在第 2 步以不带时区的
timestamp追加,被 schema 不一致拦下;在第 5 步把旧路径在提交之后才去读取,导致两个路径相同。 - 官方文档:Multi-Engine Support · pyiceberg — API · DuckDB — Iceberg extension · Spec — Primitive Types
Spark 创建表
创建 /root/ice/eng/spark.py,应用名称为 ice-eng-spark,创建 lake.eng.orders(六个列,PARTITIONED BY (days(order_ts)),'format-version' = '2'),并用一次 append() 写入 2026-03-01 到 03-07 的七个文件。
这个快照的摘要里,engine-name 会记为 spark。评分器会检查第一次提交的行数和写入它的引擎。
Python 写入同一张表
用 /root/ice/eng/py_append.py,以 pyarrow 读取 2026-03-08 的文件(amount 为 int32,order_ts 为 UTC 时区的 timestamp),并用 load_catalog("lake").load_table("eng.orders").append(…) 追加。
pyiceberg 在写入之前会把 pyarrow schema 与表的 schema 对照。如果 order_ts 是不带时区的 timestamp,就会以与 timestamptz 不匹配为由被拒绝。分区(days)的值由 pyiceberg 自己计算并写进清单。评分器会检查第二次提交是否加上了 3 月 8 日的行数,以及是否由 Spark 之外的引擎写入。
DuckDB 读取
用 /root/ice/eng/duck.py,把 ice-loc eng.orders 打印的路径传给 iceberg_scan(),求出各地区的 sum(amount),按 {"지역": 합계, …} 的格式写入 /root/ice/eng/out/duck_region.json(占位符依次为地区与合计)。
DuckDB 不经过 Catalog,而是从一个 metadata.json 文件出发,顺着清单往下走。Spark 写的文件和 pyiceberg 写的文件在同一份列表中,所以都能读到。评分器会把它与按 3 月 1–8 日原始数据算出的合计对照。
pyiceberg 按条件读取
用 /root/ice/eng/py_read.py 统计 region == 'seoul' 且 order_ts >= 2026-03-05T00:00:00+00:00 的行数,按 {"rows": 정수} 的格式写入 /root/ice/eng/out/py_count.json(占位符为整数)。
pyiceberg 的条件会先用清单中的分区值(日期)和列统计信息选取文件,再在选中的文件内过滤行。评分器会把它与按 3 月 5–8 日原始数据算出的值对照。
陈旧指针——旧路径就是旧表
用 OLD=$(ice-loc eng.orders) 记下当前路径,用 /root/ice/eng/append.py(应用 ice-eng-append,日期参数)提交 2026-03-09,然后读取 NEW=$(ice-loc eng.orders)。让 /root/ice/eng/stale.py 接收这两个路径作为参数,分别用 DuckDB 统计行数,按 {"old_path", "old_rows", "new_path", "new_rows"} 的格式写入 /root/ice/eng/out/stale.json。
metadata 文件一旦写入就不会改变。用旧路径读取,无论什么时候读,得到的都是那一刻的表——对时间旅行有用,但如果是想读取“当前”的仪表板,就会悄悄地显示陈旧的数字。评分器会检查两个路径是否在 metadata 历史中,以及行数是否与当时的 total-records 相同。
一个引擎改的名称,另一个引擎能看到
用 /root/ice/eng/rename.py(应用 ice-eng-rename)运行 ALTER TABLE lake.eng.orders RENAME COLUMN region TO area,并用 /root/ice/eng/columns.py 把 pyiceberg 的 tbl.schema() 的列名和 DuckDB 的 iceberg_scan() 结果的列名,按 {"pyiceberg": [...], "duckdb": [...]} 的格式写入 /root/ice/eng/out/rename.json。
重命名只改变 metadata 的 schema,文件里仍保留着旧名称(region)。按字段 ID 配对的引擎,会用新名称读取旧文件的值。评分器会检查两个列表中是否有 area 而没有 region。
Spark 汇总所有引擎的行
用 /root/ice/eng/daily.py,应用名称为 ice-eng-daily,把 lake.eng.orders 按日期(to_date(order_ts))分组,把订单数和 amount 合计做成 lake.eng.daily(d, orders, amount)(CREATE TABLE … AS SELECT)。
3 月 8 日的行是 pyiceberg 写的,其余是 Spark 写的。对读取的引擎来说,没有区别。日期边界按会话时区(UTC)切分。评分器会把它与按 3 月 1–9 日原始数据算出的每天的值对照。
把多个引擎接到一张表上的规则
在 /root/ice/eng/report.md 中写出 ## 한 표, 세 엔진、## 낡은 포인터、## 이름 바꾸기 三个小节。第二节以数字写入第 5 步的 old_rows 和 new_rows。
如果要给你们团队接入新引擎(例如公司内部的 BI 工具),你会先确认什么——是否通过 Catalog,读取哪种格式版本和删除文件,如何处理 timestamptz。