レイクハウスのテーブル形式 — Apache Iceberg をメタデータで理解する
Spark が作り、Python が書き、DuckDB が読む — 一つのテーブルを三つのエンジンで
目標
Sparkが作ったテーブルにpyicebergが1日分を書き足し、DuckDBとpyicebergがそのテーブルを読み、再びSparkが、すべてのエンジンが書いた行を集めて集計します。その間に、古いmetadataのパスを握って読む「古いポインター」の落とし穴と、あるエンジンのスキーマ変更がほかのエンジンにどう見えるかを確認します。
なぜ重要なのか
テーブル形式の価値は、エンジンを選ばないことにあります。バッチはSparkで、小さなロードはPythonで、アドホックな分析はDuckDBで行っても、すべて同じmetadataと同じParquetファイルを見ます。ところが、この約束には条件が付きます。すべてのエンジンが、同じカタログを経由して「現在の」metadataを探さなければならず、同じ仕様の機能(フォーマットバージョン・削除ファイル・型)をサポートしている必要があります。 metadata.jsonのパスを直接渡して読むツールは便利ですが、そのパスが指す瞬間に固定されます。誰かがそのあとにコミットすると、皆さんは古いテーブルを読んでいるのに、何の警告もありません。型にも注意が必要です。SparkのTIMESTAMPはIcebergのtimestamptzなので、タイムゾーンなしの時刻を渡すPythonコードは、スキーマ検査に引っかかります。逆に、列名の変更のようにフィールドIDで解決される変更は、すべてのエンジンにそのまま見えます。
ステップ
- /root/ice/eng/spark.py(アプリ
ice-eng-spark)でlake.eng.ordersをdays(order_ts)で分けて作成し、2026-03-01から03-07までの7ファイルを1回でコミットしてください。 - /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に、
## 한 표, 세 엔진、## 낡은 포인터、## 이름 바꾸기の3つの節を書いてください(3つの見出しは順に、韓国語で「1つのテーブル、3つのエンジン」「古いポインター」「名前の変更」を意味する語句です)。
参考
- 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で追加しようとして、スキーマの不一致で止まること、ステップ5で古いパスをコミットのあとに読んでしまい、2つのパスが同じになることです。 - 公式ドキュメント: Multi-Engine Support・pyiceberg — API・DuckDB — Iceberg extension・Spec — Primitive Types
Sparkがテーブルを作成する
/root/ice/eng/spark.pyをアプリ名ice-eng-sparkで作成し、lake.eng.orders(列は6つ、PARTITIONED BY (days(order_ts))、'format-version' = '2')を作って、2026-03-01から03-07までの7ファイルを1回のappend()で入れてください。
このスナップショットの要約には、engine-nameがsparkと残ります。採点ツールは、最初のコミットの行数と、書き込んだエンジンを確認します。
Pythonが同じテーブルに書く
/root/ice/eng/py_append.pyで、2026-03-08のファイルをpyarrowで読み込み(amountはint32、order_tsはUTCタイムゾーンのtimestamp)、load_catalog("lake").load_table("eng.orders").append(…)で追加してください。
pyicebergは、書き込む前にpyarrowスキーマをテーブルのスキーマと比較します。order_tsがタイムゾーンなしのtimestampなら、timestamptzと合わないとして拒否します。パーティション(days)の値は、pyicebergが直接計算してマニフェストに書きます。採点ツールは、2番目のコミットが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は、カタログを経由せず、metadata.json 1つのファイルから始めて、マニフェストをたどって下ります。Sparkが書いたファイルとpyicebergが書いたファイルが同じ一覧にあるので、両方読めます。採点ツールは、3月1–8日の元データから計算した合計と比較します。
pyicebergが条件付きで読む
/root/ice/eng/py_read.pyで、region == 'seoul'かつorder_ts >= 2026-03-05T00:00:00+00:00の行数を数え、/root/ice/eng/out/py_count.jsonに{"rows": 정수}(プレースホルダーは整数です)の形で書いてください。
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が2つのパスを引数に受け取り、それぞれDuckDBで行数を数えて、/root/ice/eng/out/stale.jsonに{"old_path", "old_rows", "new_path", "new_rows"}の形で書くようにしてください。
metadataファイルは、1度書かれたら変わりません。古いパスで読めば、いつ読んでもその時点のテーブルが出ます。タイムトラベルには役に立ちますが、「現在」を読むつもりだったダッシュボードなら、黙って古い数字を見せます。採点ツールは、2つのパスが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()の結果の列名を、/root/ice/eng/out/rename.jsonに{"pyiceberg": [...], "duckdb": [...]}の形で書いてください。
名前の変更はmetadataのスキーマだけを変え、ファイルは古い名前(region)をそのまま持っています。フィールドIDで対応づけるエンジンは、新しい名前で古いファイルの値を読みます。採点ツールは、2つの一覧に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日の元データから計算した日ごとの値と比較します。
複数のエンジンを1つのテーブルにつなぐときのルール
/root/ice/eng/report.mdに、## 한 표, 세 엔진、## 낡은 포인터、## 이름 바꾸기の3つの節を書いてください(3つの見出しは順に、韓国語で「1つのテーブル、3つのエンジン」「古いポインター」「名前の変更」を意味する語句です)。2つ目の節には、ステップ5のold_rowsとnew_rowsを、数字で入れてください。
皆さんのチームに新しいエンジン(例: 社内のBIツール)をつなぐなら、何を先に確認しますか。カタログを経由するか、どのフォーマットバージョン・削除ファイルを読めるか、timestamptzをどう扱うかを、確認します。