Delta テーブルのコンカレンシー制御

複数のFabricノートブック、パイプライン、または Spark ジョブが同じ Delta テーブルに同時に書き込む場合、Delta Lake はオプティミスティック コンカレンシー制御 (OCC) を使用してテーブルの一貫性を維持します。 各トランザクションはスナップショットを読み取り、新しいファイルを書き込み、その間に競合するコミットが発生しなかったことを検証します。 競合が検出された場合、トランザクションはデータを破損するのではなく、例外で失敗します。

この記事では、Fabricで同時書き込みを管理するための実用的なパターンについて説明します。 OCC プロトコルの完全な仕様については、 Delta Lake コンカレンシー制御 (オープンソースドキュメント) を参照してください。

分離レベル

すべてのデルタ テーブルでは 、シリアル化可能分離 レベルが使用されます。 Serializable は最も厳密なレベルであり、サポートされているレベルは 1 つだけです。 これにより、同時実行トランザクションの結果が一部の順次実行順序と同じになります。

Delta Lake では、論理データを変更しない操作 ( など) にも内部 OPTIMIZE レベルが使用されます。 SnapshotIsolation では同時追加チェックがスキップされるため、コンカレント挿入と競合することなく圧縮を続行できます。 SnapshotIsolation は直接構成しません。必要に応じて、Delta Lake によって自動的に適用されます。

Serializable 分離では、並行ブラインド追加(INSERT INTO)は、同じパーティションを読み取る または MERGEUPDATE

どの操作が競合するか

すべての同時書き込みが競合するわけではありません。 重要な要因は、2 つの操作が同じ基になるファイルに触れるかどうかです。

同時実行ペア 競合ですか? なぜでしょうか
2 回の INSERT(追加)操作 いいえ それぞれが既存のファイルを読み取らずに新しいファイルを追加します (ブラインドアペンド)。
INSERT + OPTIMIZE いいえ OPTIMIZE は論理データを変更しないため、 SnapshotIsolation でコミットされるため、同時追加チェックは完全にスキップされます。 追加すると、圧縮されるファイルと重複しない新しいファイルが追加されます。
2 つの UPDATEDELETE、または MERGE 操作 はい (重複するファイルの読み取りまたは変更を行う場合) それぞれがファイルを書き換えるので、2 番目のライターのスナップショットは古くなります。
OPTIMIZE + UPDATE/DELETE/MERGE はい(同じファイルに触れる場合) OPTIMIZE ファイルを削除して読み取ります ( dataChange=falseを使用)。 データ変更操作でも同じファイルが読み取られた場合は、 ConcurrentDeleteReadException が発生します。
2 つの OPTIMIZE 実行 はい(同じファイルを選択した場合) 両方とも同じファイルセットを削除して書き直そうとし、 ConcurrentDeleteDeleteExceptionをトリガーします。
INSERT + MERGE/UPDATE/DELETE はい (データ変更操作が同じパーティションを読み取る場合) Serializableでは、操作によって追加が書き込まれたパーティションが読み取られた場合、ブラインド追加は同時データ変更と競合する可能性があります。

Tip

追加専用パイプライン (INSERT INTOdf.write.mode("append")) は、競合を完全に回避する最も簡単な方法です。 ワークロードでまず追記し、後から整合を取れるなら、書き込み同士の競合を排除できます。

パーティショニングでライターを分離する

競合なしで同じテーブルに対して同時実行 DML を実行する最も一般的な方法は、ライターを区切る列で テーブルをパーティション分割 し、その列をすべての操作条件に含める方法です。 各ライターが異なるパーティションをターゲットにしている場合、操作は不整合なファイル セットに触れ、競合しません。

一般的なシナリオ: 異なる部署またはテナントの各プロセス データを複数のパイプラインで実行します。 そのディメンションでパーティション分割し、各パイプラインの MERGE をそのパーティションにピン留めします。

-- Each pipeline targets its own partition, so concurrent runs don't conflict
MERGE INTO events AS target
USING staged AS source
ON target.event_id = source.event_id
    AND target.business_unit = 'EMEA'
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

Important

パーティション列は、ソース データだけでなくマージ条件自体に含まれている必要があります。 これを使用しないと、Delta Lake は検証時に、2 つの操作が不整合なファイル セットに触れたことを判断できません。競合チェッカーは、操作をテーブル全体の読み取りとして扱います。

パーティション分割戦略の詳細については、「 差分テーブルのパーティション分割」を参照してください。

組み込み型コミット再試行

Delta Lake は、別のトランザクションが最初にコミットされたことを検出すると、コミットを自動的に再試行します。 再試行のたびに、成功したコミットを読み取り、競合チェッカーを実行し、論理競合が存在しない場合は、次の使用可能なバージョンでコミットを再試行します。 このプロセスは、コードからのアクションなしで透過的に繰り返されます。

論理競合 (たとえば、同じファイルを書き換える 2 つの操作) を自動的に解決することはできません。 再試行により、 一般的な競合例外に記載されている例外のいずれかが発生します。 ただし、同じバージョン スロットに対する 2 つのブラインド アペンド レーシングなど、多くの一時的なバージョンの競合は自動的に解決され、アプリケーションには表示されません。

一般的な競合の例外

競合が検出されると、Delta Lake によって特定の例外が発生します。 表示される例外を理解すると、根本原因を特定するのに役立ちます。

例外 どうされました
ConcurrentAppendException 別のライターが、実行していた操作で読み取っていたパーティション (またはファイル セット) にファイルを追加しました。 MERGEが、別のパイプラインから挿入を受け取るパーティションに対して実行される場合に一般的です。 Serializable 分離レベルでは、ブラインド追加でさえ、単純な INSERT 操作でもこの例外を引き起こすことがあります。
ConcurrentDeleteReadException 別のライターが、操作で読み取ったファイルを削除または書き直しました。 OPTIMIZE が、同時実行中の UPDATE または MERGE も読み取っていたファイルを圧縮する場合、または同じ行に対する2つのデータ変更操作が重複する場合によく発生します。
ConcurrentDeleteDeleteException どちらの操作も、同じファイルを削除または書き換えようとしました。 多くの場合、重複する OPTIMIZE 実行または 2 つのパイプラインで同じパーティションが同時に書き換えられたことが原因です。
ConcurrentWriteException 競合解決が実行される前に、別のトランザクションが同じテーブル バージョンにコミットしていた場合に発生する一般的な競合エラーです。たとえば、ファイルシステムから管理コミットへのアップグレード中に発生します。
MetadataChangedException テーブル スキーマまたはプロパティは、トランザクションの途中で変更されました。たとえば、同時 ALTER TABLE 書き込みやスキーマの進化書き込みなどです。
ConcurrentTransactionException 同じチェックポイントの場所を持つ 2 つの構造化ストリーミング クエリが同時にテーブルに書き込まれます。 ストリーミング ジョブを重複除去するか、個別のチェックポイント パスを使用します。
ProtocolChangedException 現在のトランザクションもプロトコルの変更を試みている間に、同時実行トランザクションがテーブル プロトコルをアップグレードまたはダウングレードしました。 テーブル機能が同時に削除された場合にも発生する可能性があります。

書き込みの競合を回避するための一般的な戦略

自動圧縮を有効にする

自動圧縮は 、書き込み操作の一部として同期的に実行されます。 同期コンパクションは、個別にスケジュールされたコンパクションジョブがデータ変更操作と重ならないようにし、それによって同時書き込み例外の発生を防ぎます。

書き込み期間外のメンテナンスをスケジュールする

OPTIMIZEVACUUM は、同時データ変更操作と競合する可能性があります。 Fabric では、テーブルのコンパクションVACUUM のためのノートブック ジョブまたはパイプライン アクティビティは、アクティビティが少ない時間帯にスケジュールします。たとえば、夜間インジェストの実行中ではなく、完了後にスケジュールします。

追加 + マージ パターンを使用する

コンカレンシーの高いインジェストの場合は、追加のみの書き込みを含む生データをステージング テーブルに配置し (競合は発生しません)、1 つの MERGE ジョブを実行してターゲット テーブルに調整します。 このパターンは、インジェストを完全に並列に保ちながら、競合しやすい操作をシリアル化します。

論理競合の再試行ロジックを追加する

組み込みのコミット再試行では、一時的なバージョンの競合が自動的に処理されますが、論理的な競合 (2 つの操作が真に重複する場合) によって例外が発生します。 Delta Lake は部分的な書き込みを生成しないため、失敗したトランザクションはアプリケーション レベルで安全に再試行できます。 論理競合が時折発生することが予想されるパイプラインの場合は、書き込みを再試行ロジックでラップします。

from delta.exceptions import ConcurrentAppendException
import time

# Retry with backoff on transient concurrent write conflicts
max_retries = 3
for attempt in range(max_retries):
    try:
        spark.sql("MERGE INTO target USING source ON ...")
        break
    except ConcurrentAppendException:
        if attempt < max_retries - 1:
            time.sleep(2 ** attempt)
        else:
            raise

適切なレイアウト戦略を選択する

液体クラスタリングパーティション分割によって 、さまざまな問題が解決されます。 Liquid クラスタリングは、読み取りパフォーマンスのためにファイル レイアウトを最適化します。 パーティション分割により、同時ライターの競合を防ぐ物理的な境界が作成されます。 ワークロードで両方が必要な場合は、ライター分離列でパーティション分割し、読み取りパフォーマンスのために各パーティション内 で Z オーダー を使用します。