レイクハウスのテーブル形式 — Apache Iceberg をメタデータで理解する
一つのテーブル、複数のエンジン — 仕様が契約でカタログが出会う場所だ
一言でいうと
Icebergのテーブルは、仕様どおりに書かれたmetadataとデータファイルなので、Sparkが作ったテーブルにPythonが書き込み、DuckDBが読み取れます。ただしこの約束は、全員が同じカタログで「現在の」metadataを探し、同じ仕様の機能と型を理解しているときにだけ守られます。
なぜ複数のエンジンなのか
バッチ変換はSparkが、小さなロードとチェックスクリプトはPythonが、アナリストのアドホックなクエリはDuckDBが、それぞれ扱いやすいです。以前はエンジンごとにテーブルを別に持ち、コピーしていました。コピーは常にずれていき、どちらが本物かで揉めることになります。Multi-Engine Supportドキュメントは、Icebergを、どの処理エンジンからでも使えるオープンスタンダードだと紹介しています。テーブルがファイル形式ではなく、仕様が定めるmetadataで定義されているからこそ、可能なことです。
同じドキュメントは、条件も示しています。SparkとFlinkは、エンジンのバージョンごとにランタイムjarが別々に出ていて、サポート一覧にないバージョンでは使えません。このコースのラボイメージが、リポジトリのほかのSparkコース(4.2)と違って、Spark 4.1.3を使う理由がそれです。Iceberg 1.11.0のランタイムは、4.1までです。
どう動くのか: 出会う場所とずれる場所
カタログが出会う場所です。ラボ環境では、Sparkとpyicebergが同じSQLiteカタログを開きます。どちらも条件付き置換でコミットするので、互いのコミットを上書きしません。DuckDBのiceberg拡張は、2つの方式を提供します。metadataを指して直接読み取る方式は、カタログが不要で読み取り専用であり、書き込みまで行うには、RESTカタログをATTACHします。
直接読み取りは、その瞬間に固定されます。metadataファイルは1度書かれたら変わらないので、iceberg_scan('…/00003-….metadata.json')は、いつ実行してもそのときのテーブルを返します。誰かがそのあとにコミットしても、警告はありません。DuckDBのドキュメントは、ファイル名で「最新」のバージョンを推測する機能がACIDを壊すおそれがあるため、デフォルトではオフになっていると記しています。現在を読むには、毎回カタログから現在のパスを探し直す必要があります。
型がずれる場所です。仕様のプリミティブ型は、タイムゾーンなしのtimestampと、UTCで保存するtimestamptzを区別します。SparkのTIMESTAMPはtimestamptzになります。そのため、Pythonでタイムゾーンなしの時刻で同じ列に書き込もうとすると、pyicebergがスキーマが合わないと拒否します(ラボイメージで確認)。時刻にUTCを付けるのが正解です。日付の境界もUTCで切られることを、覚えておいてください。
機能がずれる場所です。フォーマットバージョンは、古い読み取り側が新しい機能を正しく読めないときに上がります。仕様は、エンジンがまだ実装していない機能を避けるために、古いバージョンで書き続けてもよいと記しています。バージョン3のdeletion vectorや新しい型を使う前に、そのテーブルを読むすべてのエンジンが対応しているかを確認する必要があります。削除ファイルを知らないエンジンは、削除した行を返します。
うまく合う場所です。列はフィールドIDで選ぶので、1つのエンジンで列名を変えると、ほかのエンジンも新しい名前で古いファイルの値を読みます。パーティション値はマニフェストにあるので、誰が書いたかにかかわらず、同じ条件でスキップします。
誰が書いたのか: スナップショットの要約
エンジンが複数あると、「このコミットはどこから来たのか」を知りたくなるときが来ます。スナップショットの要約が手がかりです。ラボイメージで見ると、Sparkは要約にengine-name(spark)・engine-version・app-idを残し、pyiceberg 0.12は残しません。要約は書く側が埋めるものなので、エンジンごとに違うことも、一緒に覚えておきます。
現場での姿
ダッシュボードが昨日の数字で止まった場合です。誰かがBIツールにmetadata.jsonのパスを埋め込んでいました。テーブルは毎日コミットされるのに、そのツールは初日のテーブルを読んでいました。ツールがカタログを経由するようにするか、毎回現在のパスを探すように変えます。
Pythonのロードがスキーマエラーで止まった場合です。CSVから読んだ時刻が、タイムゾーンなしの値でした。タイムゾーンを付ければ解決します。このエラーは、ありがたく受け止めるべきです。黙って入っていたら、9時間ずれた値が混ざっていたでしょう。
新しいエンジンをつなぐ場合です。カタログをサポートしているか、どのフォーマットバージョン・削除形式を読めるか、timestamptzをどう扱うかを、先に確認します。3つのうち1つでもずれると、結果が静かに間違います。
実務で本当に大切なこと
- カタログは1つ、「現在」は毎回カタログから得ます。metadataのパスを握りしめると、古くなります。
- timestampとtimestamptzを区別します。SparkのTIMESTAMPはtimestamptzです。
- フォーマットバージョンと削除形式は、最も遅れたエンジンに合わせます。
- 名前の変更のように、IDで解決される変更は、すべてのエンジンにそのまま見えます。
次のラボですること
Sparkが日付で分けたテーブルを作って1週間分を入れ、pyicebergがタイムゾーン付きの時刻で、さらに1日分を書き込みます。DuckDBで地域別の合計を、pyicebergで条件付きの読み取りを行ったあと、古いmetadataのパスを書き留めて、Sparkにもう1日分をコミットさせ、古いパスと新しいパスをDuckDBでそれぞれ読んでみます。Sparkで列名を変更して、ほかの2つのエンジンが新しい名前で読めるかを確認し、最後に、Sparkが3つのエンジンが書いた行をすべて集めて、日別のテーブルを作ります。