データ接続と統合Iceberg tablesカタログの動作トランザクション

注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

トランザクション

Foundry の Iceberg カタログは、標準の Iceberg を拡張し、すべて成功するか、何も反映されないトランザクションセマンティクスを提供します。Foundry のビルドシステムを通じてジョブが複数の書き込みを実行する場合、すべての書き込みがまとめて成功するか、完全に破棄されます。これは、Foundry カタログのデータセットに対する既存の保証と同じです。このページでは、Foundry の Iceberg トランザクションセマンティクスの仕組み、Iceberg のデフォルトの動作との比較、およびパイプラインを作成する際の意味について説明します。

このページで説明するトランザクションの保証は、Foundry のビルドシステムを通じてジョブを実行する場合にのみ適用されます。外部クライアントから Foundry の Iceberg カタログに書き込む場合は、標準の Iceberg のトランザクション動作が適用されます。その場合のトランザクションの仕組みについては、PyIceberg API ドキュメント ↗を参照してください。

Foundry Iceberg のトランザクションセマンティクス

Foundry でビルドを実行する場合、Foundry はすべて成功するか、何も反映されない更新を提供します。Foundry は、すべての Iceberg テーブルの読み取りと書き込みを自動的に1つのトランザクションにまとめます。この動作を設定するためにユーザーが操作する必要はありません。これは、Foundry のビルドシステムを使用して Iceberg テーブルに書き込む際のデフォルトの動作です。

一方、標準の Iceberg には、複数の更新にわたるトランザクションの保証がありません。代わりに、アトミックな更新を提供し、各更新は個別に適用されます。これは楽観的同時実行を可能にする意図的な設計上の選択であり、複数の書き込み処理が同じテーブルに対して同時に動作できます。このモデルは、書き込みが1回だけのトランザクションではうまく機能しますが、1つのトランザクション内で複数の書き込みを実行するパイプラインでは、更新先が1個のテーブルでも複数のテーブルでも、正確性の問題が生じる可能性があります。Foundry の、すべて成功するか何も反映されないトランザクションモデルは、ジョブ内のすべての書き込みにわたってコミットを調整することでこの問題に対処し、部分的な更新がほかの処理から見えることはありません。

具体的には、Foundry のトランザクションモデルは次の保証を提供します。

  • すべて成功するか、何も反映されないコミット:ジョブ内のすべてのテーブル更新はまとめてコミットされます。ジョブがいずれかの時点で失敗した場合、下流のコンシューマーやほかのジョブから部分的な書き込みが見えることはありません。これは、カタログデータセットのトランザクションの動作と同じです。
  • 反復可能読み取り: ジョブ内でテーブルを複数回読み取る場合、読み取りの間にテーブルが外部から更新されても、同じデータが返されます。ジョブの実行中、パイプラインからは入力の安定した一貫性のあるビューが見えます。
  • ジョブは自身の書き込みを読み取り可能: ジョブ内で行われた書き込みは、同じジョブ内の後続の読み取りですぐに参照できます。トランザクションがコミットされるまでは、ほかのどのジョブや外部の読み取り処理からも参照できません。
  • 複数テーブルのスナップショット分離: Foundry は、トランザクションの開始時に、すべての入力テーブルの一貫したスナップショットを取得します。これにより、ジョブの実行途中に入力に対する外部からの更新の一部だけが読み取られることを防ぎ、データの来歴(入力のどのバージョンから出力のどのバージョンが生成されたか)が正しく記録されることを保証します。

例:インクリメンタル(差分処理)結合パイプライン

例を使って違いを説明します。

orders と customers という2個のテーブルから新しい行の差分だけを読み取り、それらを結合して、結果を order_summary と customer_metrics の2個のテーブルに追加するジョブがあるとします。ある日、最初のテーブル(order_summary)への書き込みを処理した後、2個目のテーブル(customer_metrics)に書き込む前に、ジョブが実行途中で失敗しました。

標準の Iceberg トランザクションセマンティクスの場合: 最初の書き込みはそのまま残ります。order_summary テーブルには新しいデータバッチが反映されますが、customer_metrics には反映されません。2個の出力は一貫性のない状態になります。その後ジョブを再試行すると、同じ入力行を再度読み取り、order_summary に再度書き込むため、前回の部分的な実行で書き込まれたデータが重複します。

Foundry Iceberg トランザクションセマンティクスの場合: 最初の書き込みはコミットされません。order_summary テーブルには新しいデータバッチが反映されません。両方のテーブルは変更されず、一貫した状態のままです。その後ジョブを再試行すると、重複することなく、新しいデータが両方のテーブルに正しく書き込まれます。

部分的なトランザクションの例。

ジョブのキューイングと楽観的同時実行

Foundry のビルドシステムは、どの時点でも、特定の出力に対して実行されるジョブが最大1件であることを保証します。同じ出力に書き込むジョブはキューに入れられ、同時ではなく順番に実行されます。そのため、実際には通常のパイプラインジョブ間で書き込みの競合は発生しませんが、パイプラインジョブと Iceberg メンテナンスジョブの間では発生する可能性があります。

一方、標準の Iceberg では楽観的同時実行を使用し、複数の書き込み処理が同じテーブルに対して同時に動作でき、コミット時に競合を検出して解決します。Foundry のトランザクションモデルは、楽観的同時実行を採用する代わりに、更新と正確性に対してより厳格なアプローチを取ります。トランザクションで競合する同時更新が発生した場合、コミットに失敗し、ビルドシステムがジョブを再試行します。更新の競合が発生した場合にはジョブを再計算する必要がありますが、正確性は常に維持されます。

競合の一般的な原因は、コンパクションなどのメンテナンスタスクです。現在、これらのタスクは通常のパイプラインジョブと同時に実行されます。ジョブの実行中にコンパクションタスクがテーブルを更新すると、ジョブのトランザクションのコミットに失敗し、ビルドシステムがジョブを再試行します。コンパクションが頻繁に行われるテーブルでは、これにより再試行が時折発生する可能性があります。