TT Lab
はじめる
学ぶ 学習パス コース

レイクハウスのテーブル形式 — Apache Iceberg をメタデータで理解する

Spark が作り、Python が書き、DuckDB が読む — 一つのテーブルを三つのエンジンで

TT Labで続きを見る

目標

Sparkが作ったテーブルにpyicebergが1日分を書き足し、DuckDBとpyicebergがそのテーブルを読み、再びSparkが、すべてのエンジンが書いた行を集めて集計します。その間に、古いmetadataのパスを握って読む「古いポインター」の落とし穴と、あるエンジンのスキーマ変更がほかのエンジンにどう見えるかを確認します。

なぜ重要なのか

テーブル形式の価値は、エンジンを選ばないことにあります。バッチはSparkで、小さなロードはPythonで、アドホックな分析はDuckDBで行っても、すべて同じmetadataと同じParquetファイルを見ます。ところが、この約束には条件が付きます。すべてのエンジンが、同じカタログを経由して「現在の」metadataを探さなければならず、同じ仕様の機能(フォーマットバージョン・削除ファイル・型)をサポートしている必要があります。 metadata.jsonのパスを直接渡して読むツールは便利ですが、そのパスが指す瞬間に固定されます。誰かがそのあとにコミットすると、皆さんは古いテーブルを読んでいるのに、何の警告もありません。型にも注意が必要です。SparkのTIMESTAMPはIcebergのtimestamptzなので、タイムゾーンなしの時刻を渡すPythonコードは、スキーマ検査に引っかかります。逆に、列名の変更のようにフィールドIDで解決される変更は、すべてのエンジンにそのまま見えます。

ステップ

  1. /root/ice/eng/spark.py(アプリice-eng-spark)でlake.eng.ordersをdays(order_ts)で分けて作成し、2026-03-01から03-07までの7ファイルを1回でコミットしてください。
  2. /root/ice/eng/py_append.py(pyiceberg)で、2026-03-08を同じテーブルに追加してください。
  3. /root/ice/eng/duck.py(DuckDB)で現在のmetadataを読み、地域別のamountの合計を、/root/ice/eng/out/duck_region.jsonに書いてください。
  4. /root/ice/eng/py_read.py(pyiceberg)で、region = 'seoul'かつorder_ts >= 2026-03-05の行数を、/root/ice/eng/out/py_count.jsonに書いてください。
  5. 現在のmetadataパスを書き留めてから、/root/ice/eng/append.py(アプリice-eng-append)で2026-03-09をコミットし、/root/ice/eng/stale.pyで古いパスと新しいパスをそれぞれDuckDBで読み、/root/ice/eng/out/stale.jsonに書いてください。
  6. /root/ice/eng/rename.py(アプリice-eng-rename)でregionをareaに変更し、/root/ice/eng/columns.pyでpyicebergとDuckDBが見る列名を、/root/ice/eng/out/rename.jsonに書いてください。
  7. /root/ice/eng/daily.py(アプリice-eng-daily)で、日ごとの注文数・金額合計のテーブルlake.eng.daily(d, orders, amount)を作成してください。
  8. /root/ice/eng/report.mdに、## 한 표, 세 엔진、## 낡은 포인터、## 이름 바꾸기の3つの節を書いてください(3つの見出しは順に、韓国語で「1つのテーブル、3つのエンジン」「古いポインター」「名前の変更」を意味する語句です)。

参考

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をどう扱うかを、確認します。