スナップショットフローから自動 CDC を作成する
create_auto_cdc_from_snapshot_flow関数は、LakeFlow Pipelinesのチェンジデータキャプチャ (CDC) 機能を使用して、データベースのスナップショットからソースデータを処理するフローを作成します。AUTO CDC FROM SNAPSHOT の仕組みを参照してください。
この関数は以前の関数apply_changes_from_snapshot()を置き換えます。2 つの関数は同じシグネチャを持ちます。Databricks では、新しい名前を使用するように更新することをお勧めします。
この操作を行うには、ターゲットストリーミングテーブルが必要です。必要なターゲットテーブルを作成するには、create_streaming_table()関数を使用できます。
構文
from pyspark import pipelines as dp
dp.create_auto_cdc_from_snapshot_flow(
target = "<target-table>",
source = Any,
keys = ["key1", "key2", "keyN"],
stored_as_scd_type = "1",
track_history_column_list = None,
track_history_except_column_list = None,
once = False
)
AUTO CDC FROM SNAPSHOT処理の場合、デフォルトの動作では、同じキーを持つ一致するレコードがターゲットに存在しない場合に新しい行を挿入します。一致するレコードが存在する場合は、行のいずれかの値が変更された場合にのみ更新されます。ターゲットには存在するがソースには存在しないキーを持つ行は削除されます。
スナップショットを使用した CDC 処理の詳細については、「The AUTO CDC APIs: Simplify チェンジデータキャプチャ with パイプライン」を参照してください。create_auto_cdc_from_snapshot_flow() 関数の使用例については、定期的なスナップショットの取り込み、履歴スナップショットの取り込み、および SCD Type 1 バックフィルの例を参照してください。
パラメーター
パラメーター | Type | 説明 |
|---|---|---|
|
| 必須。更新するテーブルの名前。 |
|
| 必須。定期的にスナップショットを作成するテーブルまたはビューの名前、または処理するスナップショット データフレーム とスナップショット バージョンを返す Python ラムダ関数。 |
|
| 必須。ソースデータ内の行を一意に識別する列または列の組み合わせ。これは、どのCDCイベントがターゲットテーブル内の特定のレコードに適用されるかを識別するために使用されます。 次のいずれかを指定できます。
|
|
| レコードを SCD タイプ 1 として保存するか、SCD タイプ 2 として保存するかを指定します。SCD タイプ 1 の場合は |
|
| ターゲット テーブル内の履歴を追跡する出力列のサブセット。追跡する列の完全なリストを指定するには、
|
|
| スナップショットフローを 1 回だけランするかどうか。 |
注意
トリガーされたパイプライン内のSCDタイプ1ターゲットの場合、1つのcreate_auto_cdc_from_snapshot_flow()フローは、PythonまたはSQLで定義された1つ以上のAUTO CDCフローとターゲットを共有できます。次の要件が適用されます。
- 各
AUTO CDCフローに一意の名前を付けます。 - すべてのフローで同じ順序で同数のキーを使用します。スナップショットフローのキー名の大文字と小文字は、
AUTO CDCのキー名と比較されます。複数のAUTO CDCフローでは、同一のキー名と大文字・小文字を使用する必要があります。 - スナップショットのバージョンと、各
AUTO CDCフローのシーケンス列にはまったく同じデータ型を使用してください。 AUTO CDC FROM SNAPSHOTフローでエクスペクテーションを定義しないでください。IGNORE NULL UPDATESを使用しないでください。Python では、create_auto_cdc_flow()フローにignore_null_updates、ignore_null_updates_column_list、またはignore_null_updates_except_column_listを設定しないでください。AUTO CDCフローを既存のAUTO CDC FROM SNAPSHOTターゲットに追加する場合、ターゲットには、予約済みのAUTO CDCシステム列と名前が競合するユーザー列を含めてはなりません。- SCDタイプ2または二時点ターゲットにはこのパターンを使用しないでください。
例については、AUTO CDC SCD Type 1 テーブルへのバックフィルの追加を参照してください。
source引数を実装する
create_auto_cdc_from_snapshot_flow() 関数には source 引数が含まれています。履歴スナップショットを処理する場合、 引数は、処理するスナップショットデータを含む とスナップショットバージョンという sourcePython2 つの値を 関数に返す ラムダ関数である必要があります。create_auto_cdc_from_snapshot_flow()Pythonデータフレーム
関数の最初の呼び出しでは、スナップショット DataFrame とバージョンを返す必要があります。最初の呼び出しで None が返された場合、パイプラインの更新は失敗します。少なくとも 1 つのスナップショットが処理された後、追加のスナップショットがないことを示すために None を返します。
以下はラムダ関数のシグネチャです。
lambda Any => Optional[(DataFrame, Any)]
- ラムダ関数への引数は、最後に処理されたスナップショット バージョンです。
- ラムダ関数の戻り値が
Noneまたは 2 つの値のタプルである: タプルの最初の値は、処理するスナップショットを含む データフレーム です。 タプルの 2 番目の値は、スナップショットの論理順序を表すスナップショット バージョンです。
ラムダ関数を実装して呼び出す例:
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Tuple[DataFrame, Optional[int]]:
if latest_snapshot_version is None:
return (spark.read.load("filename.csv"), 1)
else:
return None
create_auto_cdc_from_snapshot_flow(
# ...
source = next_snapshot_and_version,
# ...
)
The LakeFlow Pipelines ランタイムは、create_auto_cdc_from_snapshot_flow() 関数を含むパイプラインがトリガーされるたびに、次のステップを実行します。
next_snapshot_and_version関数を実行して、次のスナップショット データフレーム と対応するスナップショット バージョンを読み込みます。- 最初のエボケーションで DataFrame が返されない場合、更新は失敗します。後のエボケーションで DataFrame が返されない場合、ランが終了し、パイプラインの更新は完了としてマークされます。
- 新しいスナップショットの変更を検出し、それをターゲット テーブルに段階的に適用します。
- ステップ #1 に戻り、次のスナップショットとそのバージョンをロードします。