メインコンテンツまでスキップ

merge を使用した Delta Lake テーブルへのアップサート

MERGE SQL操作を使用して、ソーステーブル、ビュー、またはDataFrameからターゲットDelta Lakeテーブルにデータをアップサートできます。Delta Lakeは、MERGEでの挿入、更新、削除をサポートしており、高度なユースケースを容易にするためにSQL標準を超える拡張構文をサポートしています。

people10mupdatesという名前のソーステーブル、または/tmp/delta/people-10m-updatesのソースパスがあり、people10mという名前のターゲットテーブル、または/tmp/delta/people-10mのターゲットパスの新しいデータが含まれているとします。これらの新しいレコードの一部は、ターゲットデータにすでに存在している可能性があります。新しいデータをMergeするには、その人のidが既に存在する行を更新し、一致するidが存在しない新しい行を挿入します。

この例を実行できるようにターゲットテーブルとソーステーブルを作成するには、以下を実行します。ソースはターゲットと1つの id を共有しており、Merge はこれを更新します。また、1つの新しい id が追加され、Merge はこれを挿入します:

SQL
CREATE OR REPLACE TABLE main.default.people10m (
id INT, firstName STRING, middleName STRING, lastName STRING,
gender STRING, birthDate TIMESTAMP, ssn STRING, salary INT
);
INSERT INTO main.default.people10m VALUES
(1, 'Angela', 'Marie', 'Vasquez', 'F', TIMESTAMP '1980-05-12', '111-11-1111', 72000),
(2, 'Brian', 'Lee', 'Chen', 'M', TIMESTAMP '1975-09-30', '222-22-2222', 88000),
(3, 'Carla', '', 'Nguyen', 'F', TIMESTAMP '1990-01-22', '333-33-3333', 65000);

CREATE OR REPLACE TABLE main.default.people10mupdates (
id INT, firstName STRING, middleName STRING, lastName STRING,
gender STRING, birthDate TIMESTAMP, ssn STRING, salary INT
);
INSERT INTO main.default.people10mupdates VALUES
(2, 'Brian', 'Lee', 'Chen', 'M', TIMESTAMP '1975-09-30', '222-22-2222', 95000),
(4, 'Dana', 'Rae', 'Okafor', 'F', TIMESTAMP '1988-11-03', '444-44-4444', 70000);

次のクエリーを実行して、新しいデータをMergeします:

SQL
MERGE INTO main.default.people10m AS people10m
USING main.default.people10mupdates AS people10mupdates
ON people10m.id = people10mupdates.id
WHEN MATCHED THEN
UPDATE SET
id = people10mupdates.id,
firstName = people10mupdates.firstName,
middleName = people10mupdates.middleName,
lastName = people10mupdates.lastName,
gender = people10mupdates.gender,
birthDate = people10mupdates.birthDate,
ssn = people10mupdates.ssn,
salary = people10mupdates.salary
WHEN NOT MATCHED
THEN INSERT (
id,
firstName,
middleName,
lastName,
gender,
birthDate,
ssn,
salary
)
VALUES (
people10mupdates.id,
people10mupdates.firstName,
people10mupdates.middleName,
people10mupdates.lastName,
people10mupdates.gender,
people10mupdates.birthDate,
people10mupdates.ssn,
people10mupdates.salary
)
重要

ソース・テーブルの 1 つの行のみが、ターゲット・テーブルの特定の行に一致します。 Databricks Runtime 16.0 以降では、 MERGE 句と WHEN MATCHED 句と ON 句で指定された条件を評価して、重複する一致を判断します。 Databricks Runtime 15.4 LTS 以下では、 MERGE 操作では ON 句で指定された条件のみが考慮されます。

ScalaおよびPythonの構文の詳細については、 Delta Lake APIのドキュメントを参照してください。SQL構文の詳細については、 MERGE INTOを参照してください。

一致しないすべての行を merge を使用して変更する​

Databricks SQL および Databricks Runtime 12.2 LTS 以降では、 WHEN NOT MATCHED BY SOURCE 句を使用して、ソース テーブルに対応するレコードがないターゲット テーブル内のレコードを UPDATE または DELETE できます。 Databricks では、ターゲット テーブルが完全に書き換えられないように、省略可能な条件句を追加することをお勧めします。

以下のコード例は、これを削除に使用し、ターゲットテーブルをソーステーブルの内容で上書きし、ターゲットテーブル内の一致しないレコードを削除する基本的な構文を示しています。ソースの更新と削除に時間制限があるテーブルのよりスケーラブルなパターンについては、Delta Lake のテーブルをソースと増分同期するを参照してください。

Python
(targetDF.alias("target")
.merge(sourceDF.alias("source"), "source.key = target.key")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.whenNotMatchedBySourceDelete()
.execute()
)

次の例では、WHEN NOT MATCHED BY SOURCE句に条件を追加し、一致しないターゲット行で更新する値を指定します。

Python
(targetDF.alias("target")
.merge(sourceDF.alias("source"), "source.key = target.key")
.whenMatchedUpdate(
set = {"target.lastSeen": "source.timestamp"}
)
.whenNotMatchedInsert(
values = {
"target.key": "source.key",
"target.lastSeen": "source.timestamp",
"target.status": "'active'"
}
)
.whenNotMatchedBySourceUpdate(
condition="target.lastSeen >= (current_date() - INTERVAL '5' DAY)",
set = {"target.status": "'inactive'"}
)
.execute()
)

マージ操作のセマンティクス​

以下は merge プログラム操作のセマンティクスの詳細である。

  • whenMatched句とwhenNotMatched句はいくつでも指定できます。

  • whenMatched 句は、一致条件に基づいてソース行がターゲットテーブルの行と一致する場合に実行されます。これらの句には、次のセマンティクスがあります。

    • whenMatched 句には、最大で 1 つの update アクションと 1 つの delete アクションを含めることができます。 merge の update アクションは、一致したターゲット行の指定された列のみを更新します (update 操作と同様)。delete アクションは、一致したローを削除します。

    • 各whenMatched句には、省略可能な条件を設定できます。この句条件が存在する場合、句条件がtrueの場合にのみ、一致するソース行とターゲット行のペアに対してupdateまたはdeleteアクションが実行されます。

    • 複数のwhenMatched句がある場合、それらは指定された順序で評価されます。最後の句を除くすべてのwhenMatched句には条件が必要です。

    • マージ条件に一致するソース行とターゲット行のペアについて、どのwhenMatched条件もtrueと評価されない場合、ターゲット行は変更されないままになります。

    • ターゲット Delta Lake テーブルのすべての列をソースデータセットの対応する列で更新するには、whenMatched(...).updateAll()を使用します。これは次と同等です:

      Scala
      whenMatched(...).updateExpr(Map("col1" -> "source.col1", "col2" -> "source.col2", ...))

      ターゲットのDelta Lakeテーブルのすべての列に対する処理です。したがって、このアクションでは、ソーステーブルにターゲットテーブルの列と同じ列があることが前提となっています。そうでない場合、クエリは分析エラーをスローします。

注記

この動作は、自動スキーマ進化が有効になっていると変更されます。 詳細については、 自動スキーマ進化 を参照してください。

  • whenNotMatched 句は、一致条件に基づいてソース行がターゲット行と一致しない場合に実行されます。これらの句には、次のセマンティクスがあります。

    • whenNotMatched 句には、insertアクションのみを含めることができます。新しい行は、指定された列と対応する式に基づいて生成されます。ターゲットテーブルのすべての列を指定する必要はありません。指定されていないターゲット列の場合は、NULLが挿入されます。

    • 各whenNotMatched句には、省略可能な条件を設定できます。句条件が存在する場合、その行に対してその条件がtrueの場合にのみソース行が挿入されます。それ以外の場合、ソース列は無視されます。

    • 複数のwhenNotMatched句がある場合、それらは指定された順序で評価されます。最後の句を除くすべてのwhenNotMatched句には条件が必要です。

    • ターゲットDelta Lakeテーブルのすべての列を、ソースデータセットの対応する列とともに挿入するには、whenNotMatched(...).insertAll()を使用します。これは次と同等です:

      Scala
      whenNotMatched(...).insertExpr(Map("col1" -> "source.col1", "col2" -> "source.col2", ...))

      ターゲットのDelta Lakeテーブルのすべての列に対する処理です。したがって、このアクションでは、ソーステーブルにターゲットテーブルの列と同じ列があることが前提となっています。そうでない場合、クエリは分析エラーをスローします。

注記

この動作は、自動スキーマ進化が有効になっていると変更されます。 詳細については、 自動スキーマ進化 を参照してください。

  • whenNotMatchedBySource 句は、マージ条件に基づいてターゲット行がソース行と一致しない場合に実行されます。これらの句には、次のセマンティクスがあります。

    • whenNotMatchedBySource 節はdeleteとupdateのアクションを指定できる。
    • 各whenNotMatchedBySource句には、省略可能な条件を設定できます。句条件が存在する場合、ターゲット行は、その行に対してその条件がtrueである場合にのみ変更されます。それ以外の場合、ターゲット行は変更されません。
    • 複数のwhenNotMatchedBySource句がある場合、それらは指定された順序で評価されます。最後の句を除くすべてのwhenNotMatchedBySource句には条件が必要です。
    • 定義上、 whenNotMatchedBySource 句には列の値を取得するソース行がないため、ソース列を参照できません。変更する各カラムに対して、リテラルを指定するか、ターゲットカラムに対してアクション ( SET target.deleted_count = target.deleted_count + 1など) を実行できます。
重要
  • ソースデータセットで複数の行が一致し、マージによってターゲット Delta Lake テーブルの同じ行の更新が試行された場合、merge 操作が失敗する可能性があります。マージの SQL セマンティクスに従い、一致したターゲット行を更新するためにどのソース行を使用する必要があるかが不明瞭である場合、このような更新操作はあいまいとなります。ソーステーブルを前処理することで、複数一致が発生しないようにすることができます。
  • SQL VIEWにSQLMERGE操作を適用できるのは、ビューがCREATE VIEW viewName AS SELECT * FROM deltaTableとして定義されている場合のみです。

Delta Lakeテーブルへの書き込み時のデータ重複排除​

一般的な ETL のユースケースは、ログをテーブルに追加して Delta Lake テーブルに収集することです。ただし、多くの場合、ソースは重複するログレコードを生成する可能性があり、それらに対処するためにダウンストリームの重複排除ステップが必要になります。mergeを使用すると、重複するレコードの挿入を回避できます。

SQL
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId
WHEN NOT MATCHED
THEN INSERT *
注記

新しいログを含むデータセットは、それ自体内で重複排除する必要があります。マージのSQLセマンティクスにより、新しいデータをテーブル内の既存のデータと照合して重複を排除しますが、新しいデータセット内に重複データがある場合は、そのデータが挿入されます。したがって、テーブルにマージする前に、新しいデータの重複を排除してください。

重複するレコードが数日間しか表示されないことがわかっている場合は、テーブルを日付で分割し、照合するターゲットテーブルの日付範囲を指定することで、クエリをさらに最適化できます。

SQL
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS
WHEN NOT MATCHED AND newDedupedLogs.date > current_date() - INTERVAL 7 DAYS
THEN INSERT *

これは、テーブル全体ではなく、過去7日間のログのみで重複を検索するため、前のコマンドよりも効率的です。さらに、この挿入専用マージを構造化ストリーミングと使用して、ログの継続的な重複排除を実行できます。

  • ストリーミングクエリでは、foreachBatch でマージ操作を使用することで、重複排除しながらストリーミングデータを Delta Lake テーブルに継続的に書き込むことができます。foreachBatch の詳細については、以下のストリーミングの例を参照してください。
  • 別のストリーミングクエリーでは、このDelta Lakeテーブルから重複排除されたデータを継続的に読み込むことができます。挿入のみのマージは、Delta Lake テーブルに新しいデータを追加するだけであるため、このようなことが可能です。

Delta Lakeによるゆっくり変化するデータ(SCD)とチェンジデータキャプチャ(CDC)​

Lakeflow pipelines は、SCD Type 1 および Type 2 の追跡と適用をネイティブでサポートしています。CDC フィードの処理時に順不同のレコードが正しく処理されるように、Lakeflow pipelines で AUTO CDC ... INTO を使用してください。「AUTO CDC APIs: パイプラインによるチェンジデータキャプチャの簡素化」を参照してください。

Delta Lake テーブルをソースと増分同期する​

Databricks SQL および Databricks Runtime 12.2 LTS 以降では、 WHEN NOT MATCHED BY SOURCE を使用して任意の条件を作成し、テーブルの一部をアトミックに削除および置換できます。 これは、最初のデータ入力後数日間レコードが変更または削除される可能性があるが、最終的には最終的な状態に落ち着くソース テーブルがある場合に特に便利です。

次のクエリは、このパターンを使用して、ソースから5日分のレコードを選択し、ターゲットの一致するレコードを更新し、ソースからターゲットに新しいレコードを挿入し、ターゲットの過去5日間の一致しないレコードをすべて削除することを示しています。

SQL
MERGE INTO target AS t
USING (SELECT * FROM source WHERE created_at >= (current_date() - INTERVAL '5' DAY)) AS s
ON t.key = s.key
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
WHEN NOT MATCHED BY SOURCE AND created_at >= (current_date() - INTERVAL '5' DAY) THEN DELETE

ソーステーブルとターゲットテーブルに同じブールフィルタを提供することにより、削除を含む変更をソーステーブルからターゲットテーブルに動的に伝播することができます。

注記

このパターンは条件句なしで使用できますが、ターゲットテーブルを完全に書き換えることになり、コストがかかる可能性があります。