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

For eachタスクを使用して、複数のテーブルを増分でコピーします。

多数のソーステーブルからUnity Catalogテーブルにスケジュールに基づいてデータをコピーする必要がある場合、すべての実行で全行をコピーすると、処理が遅くコストがかかります。各テーブルの最後に処理された行を追跡するには、 ウォーターマーク を使用し、各実行で新しい行のみをコピーします。

このチュートリアルでは、メタデータ駆動型ジョブを構築する方法を示します。

  • ソーステーブルのリストとそのウォーターマークの状態をDeltaコントロールテーブルに格納します。
  • For eachタスクを使用して、各テーブルを並列で処理します。
  • 最終成功実行以降に追加された行のみをコピーします
  • コピーが正常に完了するたびにウォーターマークを更新します。

仕組み

ジョブでは、3つのタスクタイプが順番に連携して使用されます。

タスク

Type

機能

read_watermarks

SQL

ウォーターマークコントロールテーブルを読み取り、ソーステーブルごとに1行を返します。

copy_tables

For each

{{tasks.read_watermarks.output.rows}}を反復処理し、ソーステーブルごとにネストされたタスクを1回実行します

copy_incremental (ネスト)

ノートブック

最後のウォーターマーク以降に追加された行を読み取り、それらをターゲットテーブルに書き込み、ウォーターマークを進めます。

タスク

Type

機能

read_watermarks

SQL

ウォーターマークコントロールテーブルを読み取り、ソーステーブルごとに1行を返します。

copy_tables

For each

{{tasks.read_watermarks.output.rows}}を反復処理し、ソーステーブルごとにネストされたタスクを1回実行します

copy_incremental (ネスト)

ノートブック

最後のウォーターマーク以降に追加された行を読み取り、それらをターゲットテーブルに書き込み、ウォーターマークを進めます。

SQLタスクの出力 —行オブジェクトのJSON配列— は、For each タスクの**入力**フィールドに{{tasks.read_watermarks.output.rows}}を使用して流れ込みます。ネストされたノートブックは、各イテレーションでsource_tabletarget_tablewatermark_column、およびlast_watermarkを受け取ります。

前提条件

  • ジョブとノートブックを作成する権限を持つDatabricksワークスペース
  • Unity Catalog でスキーマとテーブルを作成する権限
  • SQLタスクを実行するためのSQLウェアハウス

このチュートリアルでは、samples.wanderbricks サンプルデータセットから読み取り、各コードブロックの先頭にある catalog および schema 変数によって指定された名前のスキーマに書き込みます。これらの変数は default に main.example_output です。別の場所に書き込むには、すべてのブロックで両方の値を一貫して変更してください。この例では、スキーマが存在しない場合に作成します。

ステップ 1: ウォーターマーク制御テーブルを作成します

ウォーターマークコントロールテーブルは、処理するテーブルを決定し、各テーブルがどれだけコピーされたかを示す信頼できるソースです。各行は1つのソーステーブルを表します。

次の SQL を実行してコントロールテーブルを作成し、2 つのソーステーブルを登録します。この例では samples.wanderbricks.userssamples.wanderbricks.properties を登録し、それぞれを独自のスキーマ内のターゲットテーブルにコピーします。catalog 変数と schema 変数は、コントロールテーブルの場所とターゲットテーブル名の両方を設定します:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

CREATE SCHEMA IF NOT EXISTS IDENTIFIER(catalog || '.' || schema);

CREATE OR REPLACE TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') (
source_table STRING NOT NULL,
target_table STRING NOT NULL,
watermark_column STRING NOT NULL,
last_watermark TIMESTAMP NOT NULL
);

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks') VALUES
('samples.wanderbricks.users', catalog || '.' || schema || '.users', 'created_at', '1970-01-01'),
('samples.wanderbricks.properties', catalog || '.' || schema || '.properties', 'created_at', '1970-01-01');

両方のソーステーブルが、created_at をウォーターマーク列として使用します。これは、新しい行が到着するたびに増加し続ける挿入時刻のTimestampです。最初のランで last_watermark1970-01-01 に設定すると、ノートブックは既存のすべての行をコピーします。これは初期フルロードとして機能します。後続のランでは、前回のラン以降に追加された行のみがコピーされます。

注記

コントロールテーブルを再作成すると、各 last_watermark1970-01-01 に Reset されますが、ターゲットテーブルはクリアされません。このステップを再実行してからジョブを再度実行すると、ノートブックは両方のソースをコピーされていないものとして扱い、すべての履歴行を 2 回目に追加します。クリーンな状態でやり直すには、コントロールテーブルを再作成する際にターゲットテーブルを削除してください。

ステップ2: コピーノートブックを作成します

ノートブックは、テーブルのイテレーションごとに1回実行されます。ウォーターマークを読み取り、ソースをフィルターし、ターゲットに書き込み、ウォーターマークを進めます。

/Workspace/Users/<username>/copy_incremental などのパスにノートブックを作成し、次のコードを追加します。Widget defaults let you ラン and test the ノートブック directly.For each タスク内で実行されると、ジョブはコントロールテーブルの場所を示す catalogschema を含め、各反復の値でそれらを上書きします。

このコードは、前回のウォーターマーク以降に追加された行のみを読み取り、ターゲットテーブルに追加します。ターゲットテーブルが存在しない場合は作成されます。次に、書き込んだばかりの行のハイウォーターマークを計算し、次のランがそこから開始されるようにコントロールテーブルを進めます:

Python
from pyspark.sql.functions import max as spark_max

# Widget defaults let you run the notebook directly; the For each task overrides them per iteration
dbutils.widgets.text("catalog", "main", "Catalog")
dbutils.widgets.text("schema", "example_output", "Schema")
dbutils.widgets.text("source_table", "samples.wanderbricks.users", "Source table")
dbutils.widgets.text("target_table", "main.example_output.users", "Target table")
dbutils.widgets.text("watermark_column", "created_at", "Watermark column")
dbutils.widgets.text("last_watermark", "1970-01-01", "Last watermark")

catalog = dbutils.widgets.get("catalog")
schema = dbutils.widgets.get("schema")
source_table = dbutils.widgets.get("source_table")
target_table = dbutils.widgets.get("target_table")
watermark_column = dbutils.widgets.get("watermark_column")
last_watermark = dbutils.widgets.get("last_watermark")

# Read only rows newer than the last watermark. A strict > can skip rows that share the
# stored high-water timestamp; for insert-only sources with distinct timestamps this is safe.
new_rows = spark.table(source_table).filter(f"{watermark_column} > '{last_watermark}'")

row_count = new_rows.count()
print(f"Copying {row_count} new rows from {source_table}")

if row_count > 0:
# Append the new rows, creating the target table on the first run
new_rows.write.format("delta").mode("append").saveAsTable(target_table)

# Compute the high-water mark from the rows just written
new_watermark = new_rows.agg(spark_max(watermark_column)).collect()[0][0]

# Advance the control table so the next run starts from here
spark.sql(f"""
UPDATE {catalog}.{schema}.watermarks
SET last_watermark = CAST('{new_watermark}' AS TIMESTAMP)
WHERE source_table = '{source_table}'
""")

print(f"Watermark for {source_table} advanced to {new_watermark}")
else:
print(f"No new rows for {source_table}, watermark unchanged")
注記

このノートブックは append モードを使用します。これは、samples.wanderbricks.userssamples.wanderbricks.properties と同様に、ソースに挿入のみが含まれる場合に適しています。ソースに更新が含まれている場合は、更新Timestampでウォーターマークを設定し、write.mode("append") の代わりに MERGE ステートメントを使用してターゲットテーブルに行をアップサートします。マージ構文については、「Merge を使用した Delta Lake テーブルへのアップサート」を参照してください。

ステップ3: ジョブを作成する

Databricksワークスペースで、サイドバーの ワークフロー をクリックし、次に ジョブ をクリックします。ジョブにIncremental table copyなどの名前を付けます。

ステップ 4: ウォーターマークルックアップタスクを構成します

SQL タスクはコントロールテーブルを読み取り、その結果を For each タスクで使用できるようにします。クエリーはコントロールテーブルを特定する catalog 変数と schema 変数を宣言するため、タスクのインライン SQL フィールドではなく、マルチステートメント SQL ファイルとして実行する必要があります。

  1. ワークスペース内に /Workspace/Users/<username>/read_watermarks.sql などの SQL ファイルを作成し、次の内容を記述します。ステップ1で使用したのと同じ値を catalogschema に設定します:

    SQL
    DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
    DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

    SELECT source_table, target_table, watermark_column, last_watermark
    FROM IDENTIFIER(catalog || '.' || schema || '.watermarks');
  2. ジョブで、 「タスクを追加」 をクリックします。

  3. [タスク名]read_watermarksに設定します。

  4. TypeSQL に設定し、 SQL タスクFile に設定します。

  5. 「パス」 に作成した SQL ファイルを設定します。

  6. **SQLウェアハウス**をワークスペース内のウェアハウスに設定します。

  7. タスクを作成 」をクリックします。

このタスクが実行されると、Databricks は結果を JSON 配列として tasks.read_watermarks.output.rows にキャプチャします。初期フルロード後、各 last_watermark はそのソースからコピーされた最新の行を反映します:

JSON
[
{
"source_table": "samples.wanderbricks.users",
"target_table": "main.example_output.users",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T23:05:18.000Z"
},
{
"source_table": "samples.wanderbricks.properties",
"target_table": "main.example_output.properties",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T00:00:00.000Z"
}
]

ステップ 5: For eachタスクを構成します

For eachタスクは SQL 出力を読み取り、ソーステーブルごとに1つのネストされたタスク実行を起動します。

  1. **タスクを追加**をクリックし、**依存元**をread_watermarksに設定します。

  2. [タスク名]copy_tablesに設定します。

  3. TypeFor each に設定します。

  4. [入力] フィールドに次を入力します:


    {{tasks.read_watermarks.output.rows}}
  5. [同時実行]2に設定すると、一度に2つのテーブルをコピーできます。ウェアハウスがより高い並列処理をサポートできる場合は、この値を増やしてください。

  6. ネストされたタスクを構成するには、[ ループするタスクを追加 ] をクリックします。

  7. [タスク名]copy_incrementalに設定します。

  8. Set Type to ノートブック .

  9. [パス] をステップ2で作成したノートブックのパスに設定します。

  10. [パラメーター] をクリックし、次に [追加] をクリックして、以下の各パラメーターを追加します。

キー

Value

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

キー

Value

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

ノートブックが SQL タスクで読み取られるコントロールテーブルを進められるように、catalogschema をステップ 1 で使用したのと同じ値に設定します。各 {{input.<key>}} 参照は、現在のイテレーションの行から対応するフィールドに解決されます。 11. 「 タスクを作成 」をクリックします。

ステップ 6: ジョブを実行し、確認します

  1. 今すぐ実行 をクリックしてジョブをトリガーします。
  2. ジョブ実行ページで、copy_tablesノードをクリックしてFor eachタスクを展開します。
  3. 実行ページには、イテレーションのテーブルが表示されます。各ソーステーブルにつき1行で、それぞれのステータス、開始時刻、期間が表示されます。
  4. 任意のイテレーションをクリックして、ノートブックの出力を表示し、行数とウォーターマークの更新を確認します。

ウォーターマークが進んだことを確認するには、ジョブ完了後に以下のクエリーを実行します。前のステップで使用したのと同じ値を catalogschema に設定します:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

SELECT source_table, last_watermark
FROM IDENTIFIER(catalog || '.' || schema || '.watermarks');

last_watermarkの値に、最後にコピーされた行のタイムスタンプが反映されます。値が1970-01-01のままである場合、ソーステーブルにフィルターに一致する行が含まれていなかったか、コピーのタスクでエラーが発生しています。詳細については、タスクの実行出力を確認してください。

パターンを拡張します。

各スニペットは、前のステップで使用された同じ catalog 変数と schema 変数を宣言します。それらをコントロールテーブルの場所を示す値に設定します。

新しいソーステーブルの追加 : コントロールテーブルに行を挿入します。次のジョブランが自動的にそれを取得し、1970-01-01 からのフルロードを開始します。このスニペットは、以下の Pause a table から active 列を既に追加済みであることを前提としているため、activetrue に設定します。2 つの拡張機能はどちらの順序でも適用できます。まだ列を追加していない場合は、最後の値とその列を削除します:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks')
(source_table, target_table, watermark_column, last_watermark, active)
VALUES
('samples.wanderbricks.hosts', catalog || '.' || schema || '.hosts', 'joined_at', '1970-01-01', TRUE);

テーブルを停止する : active 列を追加し、既存の行に対して true にバックフィルしてから、SQLファイルタスクでその列に基づいてフィルタリングします。Delta では、列の追加とその値の設定を別々のステートメントで行う必要があります:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

ALTER TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') ADD COLUMN active BOOLEAN;

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks') SET active = TRUE;

次に、ジョブが停止されたテーブルをスキップするように、read_watermarks.sql ファイル内の SELECTWHERE active = TRUE を追加します。

**テーブルをバックフィルします**:そのウォーターマークをリセットし、特定のポイントから再コピーするために

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks')
SET last_watermark = '2025-01-01'
WHERE source_table = 'samples.wanderbricks.users';

その他のリソース