レイクハウスのテーブル形式 — Apache Iceberg をメタデータで理解する
小さなコミット 30 回分を片付ける — タグ、圧縮、失効、孤立ファイルの順に
目標
小さなバッチ30個をコミット30回で入れて、ストリーミングのロードが作る小さなファイルとスナップショットの山を再現したあと、戻る時点にタグを付け、rewrite_data_filesでコンパクションし、expire_snapshotsで古いスナップショットとそのファイルを削除し、remove_orphan_filesで誰も指していないファイルを片づけます。ステップごとに、ファイル数とスナップショット数がどう変わるかを記録します。
なぜ重要なのか
Icebergは、何も自分では削除しません。コミットのたびに新しいファイルと新しいmetadataがたまり、コンパクションしても古いファイルは古いスナップショットが指しているのでそのまま残ります。そのため、整理をしていないテーブルは読み取りが遅くなり(小さなファイルやマニフェストが多いため)、ストレージは大きくなり続けます。 整理は、3つの別の仕事です。コンパクションは、小さなファイルを大きなファイルに書き直して、現在のスナップショットを速くします。期限切れ処理は、古いスナップショットと、それだけが指していたファイルを削除します。タイムトラベルで行ける過去が、その分だけ減ります。孤立ファイルの整理は、どのスナップショットも指したことがないファイル(失敗したジョブの痕跡)を削除します。現在使われている最中のファイルまで消さないように、十分な時間的余裕を置きます。 3つとも元に戻せません。そのため、順序が重要です。戻る必要がある時点があるなら、期限切れ処理の前にタグを付ける必要があり、孤立ファイルの整理の時間的余裕は縮めません。
ステップ
- /root/ice/mnt/trickle.py(アプリ
ice-mnt-trickle)でlake.mnt.ordersをPARTITIONED BY (days(order_ts))・format-version 2で作成し、/data/ice/batches/batch-001.csvからbatch-030.csvまでを、バッチごとに1コミットずつ入れてください。 - /root/ice/mnt/before.pyで、現在のスナップショット数・データファイル数・平均ファイルサイズを、/root/ice/mnt/out/before.jsonに書いてください。
- /root/ice/mnt/tag.py(アプリ
ice-mnt-tag)で、10番目のコミットのスナップショットにタグbatch10をRETAIN 30 DAYSで付けてください。 - /root/ice/mnt/compact.py(アプリ
ice-mnt-compact)でrewrite_data_filesを呼んで、結果を、/root/ice/mnt/out/compact.jsonに書いてください。 - /root/ice/mnt/expire.py(アプリ
ice-mnt-expire)で、現在より古いスナップショットをretain_last => 1で期限切れにしてください。 /root/ice/warehouse/mnt/orders/dataに、古い孤立ファイル(stray-old.parquet、更新時刻は4日前)と新しいファイル(stray-new.parquet)を置き、/root/ice/mnt/orphans.py(アプリice-mnt-orphans)でremove_orphan_filesをまずdry_runで実行して、/root/ice/mnt/out/orphans_dry.txtに書いたあと、実際に実行してください。- /root/ice/mnt/after.pyで、整理後のスナップショット数・データファイル数・ディスク上のParquetファイル数を、/root/ice/mnt/out/after.jsonに書いてください。
- /root/ice/mnt/report.mdに、
## 압축、## 만료와 태그、## 고아 파일の3つの節を書いてください(3つの見出しは順に、韓国語で「コンパクション」「期限切れ処理とタグ」「孤立ファイル」を意味する語句です)。
参考
- プロシージャは、
spark.sql("CALL lake.system.<이름>(table => 'lake.mnt.orders', …)")(プレースホルダーはプロシージャ名です)で呼びます。引数には式ではなく値(文字列・TIMESTAMPリテラル)を渡します。 remove_orphan_filesは、older_thanを現在から24時間より近く指定すると拒否します。書き込み中のファイル(まだコミット前なので孤立ファイルのように見えるファイル)を消す事故を防ぐためです。指定しなければ、3日前がデフォルトです。- よくある間違いは、タグを付ける前に期限切れ処理を動かして、10番目のスナップショットが消えることです。元に戻す方法はありません。そうなってしまったら、
DROP TABLE lake.mnt.orders PURGEを実行してから、ステップ1からやり直してください。 - 公式ドキュメント: Maintenance・Spark Procedures — rewrite_data_files · expire_snapshots · remove_orphan_files・Branching and Tagging — retention
小さなバッチ30個、コミット30回
/root/ice/mnt/trickle.pyをアプリ名ice-mnt-trickleで作成し、lake.mnt.orders(列は6つ、PARTITIONED BY (days(order_ts))、'format-version' = '2')を作って、/data/ice/batches/batch-001.csvからbatch-030.csvまで、順に1つずつappend()してください。
バッチ1つがコミット1つ、スナップショット1つで、日付パーティションごとに小さなファイルが1つずつできます。採点ツールは、30回目のコミット直後のmetadataで、30個のスナップショットがすべてそのバッチの行数を加えたappendであるかを確認します。
整理前の数字
/root/ice/mnt/before.py(pyiceberg)で、スナップショット数・データファイル数(total-data-files)・平均ファイルサイズ(total-files-size÷ファイル数、整数除算)を、/root/ice/mnt/out/before.jsonに{"snapshots", "data_files", "avg_file_bytes"}の形で書いてください。
現在のスナップショットの要約には、テーブル全体の累計値(total-*)が入っているので、マニフェストをすべて読まなくても、ファイル数とサイズがわかります。Parquetファイル1つが数KBなら、読むときにファイルを開くコストが、データを読むコストより大きくなります。
削除する前に名前を
/root/ice/mnt/tag.pyをアプリ名ice-mnt-tagで作成し、lake.mnt.orders.snapshotsをcommitted_at順に読んで、10番目のスナップショットにCREATE TAG batch10 AS OF VERSION <ID> RETAIN 30 DAYSを付けてください。
期限切れ処理は、older_thanより古いスナップショットを削除しますが、タグやブランチが指すスナップショットは、そのラベル(タグ・ブランチ)の保持期間が残っている間は削除しません。そのため、タグは期限切れ処理より先に付ける必要があります。採点ツールは、タグが30回目のコミット直後のmetadataの10番目のスナップショットを指しているかを確認します。
コンパクション: 小さなファイルを書き直す
/root/ice/mnt/compact.pyをアプリ名ice-mnt-compactで作成し、CALL lake.system.rewrite_data_files(table => 'lake.mnt.orders', options => map('min-input-files', '2'))を呼んで、結果の行のrewritten_data_files_count・added_data_files_countを、/root/ice/mnt/out/compact.jsonに{"rewritten", "added"}の形で書いてください。
コンパクションは、同じパーティションの小さなファイルを読んで大きなファイルに書き直し、「replace」スナップショット1つでコミットします。行は1つも変わりません。古い小さなファイルは、一覧から外れるだけでディスクには残ります。古いスナップショット(とタグ)がまだ指しているからです。採点ツールは、replaceコミットの要約と、皆さんの2つの値を比較します。
期限切れ処理: 古いスナップショットとそのファイルが削除される
/root/ice/mnt/expire.pyをアプリ名ice-mnt-expireで作成し、現在時刻を読んで、CALL lake.system.expire_snapshots(table => 'lake.mnt.orders', older_than => TIMESTAMP '<지금>', retain_last => 1)(プレースホルダーは現在時刻です)を呼んでください。
期限切れ処理は、スナップショットをmetadataから削除し、残ったスナップショットのどれも指していないファイルを、ディスクから削除します。現在のmainとタグbatch10のスナップショットだけが残り、タグが指す10番目のスナップショットの小さなファイルは、削除されません。採点ツールは、残ったスナップショットとディスク上のファイルを確認します。
孤立ファイルの整理: 時間的な余裕を置いて
/root/ice/warehouse/mnt/orders/dataに、既存のデータファイルを1つコピーして、stray-old.parquet(更新時刻はtouch -d '4 days ago')とstray-new.parquet(現在)を作ってください。そのあと、/root/ice/mnt/orphans.pyをアプリ名ice-mnt-orphansで作成し、remove_orphan_files(table => 'lake.mnt.orders', dry_run => true)の結果のorphan_file_locationを、/root/ice/mnt/out/orphans_dry.txtに1行に1つずつ書き、続けて、dry_runなしでもう一度呼んでください。
2つのファイルは、どちらもどのスナップショットも指していない孤立ファイルです。ところが、作ったばかりのファイルは、今まさに誰かが書いているコミットのファイルかもしれません。そのため、older_thanより新しいファイルには手を付けません(デフォルトは3日)。採点ツールは、dry_runの一覧に古いものだけがあったか、実際に古いものだけが削除されたかを確認します。
整理後の数字
/root/ice/mnt/after.py(pyiceberg)で、残ったスナップショット数・現在のスナップショットのデータファイル数・/root/ice/warehouse/mnt/orders/dataの下のParquetファイル数を、/root/ice/mnt/out/after.jsonに{"snapshots", "data_files", "files_on_disk"}の形で書いてください。
ディスク上のファイル数は、現在のスナップショットのファイル数より多くなります。タグbatch10が守る小さなファイルと、孤立ファイルの整理がわざと残した新しいファイルがあるからです。その差を説明できれば、整理処理を理解したことになります。
整理処理を運用手順にする
/root/ice/mnt/report.mdに、## 압축、## 만료와 태그、## 고아 파일の3つの節を書いてください(3つの見出しは順に、韓国語で「コンパクション」「期限切れ処理とタグ」「孤立ファイル」を意味する語句です)。最初の節には、コンパクション前のデータファイル数(ステップ2)と、コンパクション後の現在のスナップショットのデータファイル数(ステップ7)を、数字で入れてください。
この3つを毎日動くジョブにするなら、どんな順序・周期・基準(older_than、retain_last、孤立ファイルの時間的余裕)にしますか。ストリーミングで1分ごとにコミットするテーブルなら、何が変わるかも書いてみてください。