REPLACE USING フローを使用した部分的なスナップショット置換

Important

この機能は ベータ版です

REPLACE USING フローは、ターゲットテーブルをストリーミングソースと同期させ、指定されたキーカラムに一致するすべての行を置き換え、他のデータは変更しません。

SEQUENCE BY列は更新の順序を割り当て、更新が順不同になっても結果が正しいようにします。 各キーに対して、最も高い列が勝ち、低い列がターゲット内のより高い列を上書きすることはありません。 同じキーと同じシーケンスを共有する行は置き換えではなく付加されます。

REPLACE USING の仕組み

2つの地域ごとにクリックイベントとコンバージョンイベントを保持するイベントテーブルを考えてみましょう。 seq順に並んでいます。

region_id device_type event_type seq
1 iOS click 1
1 Android 変換 1
2 iOS click 1
2 デスクトップ click 1

REPLACE USING (region_id) SEQUENCE BY seqフローはリージョン1と3の更新を受け取ります。 リージョン2には更新なし:

region_id device_type event_type seq
1 iOS click 2
1 Android 変換 2
1 デスクトップ click 2
3 iOS click 1
3 デスクトップ click 2

ターゲットは次のようになります:

region_id device_type event_type seq 結果
1 iOS click 2 置き換えたのは、序列2が序列1より大きいためです
1 Android 変換 2 置き換えたのは、序列2が序列1より大きいためです
1 デスクトップ click 2 置き換えたのは、序列2が序列1より大きいためです
2 iOS click 1 このアップデートにはキーがないため、手つかずのままです
2 デスクトップ click 1 このアップデートにはキーがないため、手つかずのままです
3 デスクトップ click 2 追加。 領域3のシーケンス1行は追加されません。なぜならキーの最高位のシーケンスのみが適用されるからです。

Requirements

REPLACE USING フローには以下の要件があります:

  • REPLACE USING フローは、Databricks Runtime 18.2以降、クラシックまたはサーバーレスコンピュート上で実行されます。 DatabricksはUnity Catalogを推奨しています。
  • ソースはストリーミング ソースである必要があります。 REPLACE USING はストリーミング以外のソースを拒否します。
  • 少なくとも1つのキーカラムと、正確に1つの SEQUENCE BY カラムを指定する必要があります。

REPLACE USING フローを使うタイミング

レイクフローパイプラインは既存の列を上書きする3つの流れを提供します。 ソースの見た目や、置き換えるべき行の識別方法に基づいて選びましょう:

  • ソースが列ごとにキー化された部分的なスナップショットの連続であれば、REPLACE USING を使いましょう。 REPLACE USING は、受信データと一致するデータのみを上書きし、それ以外のデータはそのまま残します。 主キーは必要ありません。
  • ソースが明示的な挿入更新削除操作を伴う変更データキャプチャ(CDC)フィードである場合や、ゆっくりと次元を変化させる(SCD)タイプ2履歴が必要な場合はAUTO CDCを使用してください。 AUTO CDCも真のプライマリキーを必要とします。 「AUTO CDC API: パイプラインを使用して変更データ キャプチャを簡略化する」を参照してください。
  • ソースがスナップショットで、述語によって選ばれたターゲットテーブルの範囲(例えば過去7日間)をバッチ操作として再計算・上書きしたい場合、REPLACE WHEREを使いましょう。 主キーは必要ありません。 REPLACE WHERE フローを使用したバッチ処理を参照してください。

REPLACE USING フローを作成してください

SQLまたはPythonのいずれかでREPLACE USING flowsを定義してください。

SQL

FLOW REPLACE USING と同じ行で CREATE STREAMING TABLE 句を使用します。

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

または、長い形式の CREATE FLOW 構文を使用します。

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME は SQL で必要です。 位置ではなく列名で列を照合します。

Python

@dp.table でテーブルとフローをまとめて宣言します:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

あるいは、@dp.replace_flow を使用して既存のストリーミング テーブルを対象にします:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using はキーカラムのリストです。 sequence_by は列名または Column 式であり、 replace_using が設定されるたびに必要です。

順序付けと順序不同のデータ

SEQUENCE BY列は結果を更新の到着順に依存せずに作成します。 行は、そのキーの列がすでに保存されている列より大きい場合にのみキーに適用されるため、現在の値より古い遅い行や再生された行は無視されます。 アップデートに存在しないキーはそのまま残されます。

これらの手順に従い、代替が予測可能な動作をします:

演習 理由
タイムスタンプ、バージョン番号、ログオフセットなど、キーのバージョンごとに厳密に増加するシーケンスを使いましょう。 同じキーで同じシーケンスを持つ2行が両方とも保持されるため、そのキーに重複する行が発生します。
null でないシーケンスを使用してください。 ヌルシーケンスは定義されていない動作を引き起こすことがあります。

Expectations

期待を支えるフローの使用を置き換えてください。 warn fail他のフローと同様に動作します:warnは行を繰り返し違反し、違反を記録し、failは更新を停止します。 「パイプラインの期待値を使用してデータ品質を管理する」を参照してください。

drop期待値は、違反している行を、ソースがその行を一度も生成しなかったかのように扱います。 削除された行はターゲットテーブル内の一致するキーを置き換えたり削除したり修正したりしません。

  • ドロップは重複除去の前に行われるため、フローはキーの最新の有効バージョンを保持します。
  • キーの入ってくるすべての行が削除された場合、そのキーの既存の行はそのまま残ります。
  • 破棄された行はシーケンスの下限を設定しないため、後続の有効な更新は、そのシーケンスが破棄された行のシーケンスより低くても適用されます。

制限事項

REPLACE USING フローには以下の制限があります:

  • REPLACE USING はターゲットテーブルごとに単一のフローをサポートしています。 同じターゲット上でREPLACE USING と他のフロータイプを組み合わせることはサポートされていません。
  • ターゲット テーブルは、パイプライン内に作成する必要があります。
  • ソースはストリーミング ソースである必要があります。
  • 少なくとも1つのキーカラムと SEQUENCE BY カラムを指定する必要があります。 キー列は重複できず、各キー列のタイプはソート可能でなければなりません。 整数、文字列、日付などの原子型は鍵になれますが、 MAPVARIANT は鍵にはなりません。
  • スタンドアロンのストリーミングテーブルについては、構文の違いについては 「部分スナップショット置換をREPLACE USING flowsに適用 」を参照してください。

例示

以下の例は、すべてのUnityカタログ対応ワークスペースで利用可能な予約状態変更のサンプル表である samples.wanderbricks.booking_updatesから引用しています。 各予約は変更ごとに1回ずつ現れるため、 booking_id 新しい booking_update_idで繰り返されます。 Wanderbricksデータセットを参照してください。

例1:各キーの最新記録を保持する

各予約の現状のみを残してください。 フローは booking_id をキーとし、booking_update_id で順序付けされるため、予約の最新の更新がそれ以前の更新に置き換わります。 ソースが明示的な挿入、更新、削除操作を持つ変更フィードの場合は、AUTO CDCを使いましょう。

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

この例はbooking_update_idタイムスタンプではなくupdated_atでシーケンスされます。なぜなら、同じ予約の複数の更新がタイムスタンプを共有する可能性があるからです。 順序に一致する行は置き換えではなく追加されるため、予約用に複数行が残ることになります。

例2:複数の列にキーを置く

レコードが複数の列の組み合わせで識別された場合は、それらをすべて REPLACE USINGでリストアップします。 ここでは各予約が (property_id, booking_id)で識別されるため、フローは物件ごとの予約の現状を保持します。 キー列がnullでできる場合、REPLACE USING は行をスキップするのではなくnullからnullにマッチングします。

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

例3:無効なレコードを期待値でドロップする

悪い列をターゲットから遠ざけるという期待も加えましょう。 削除された行は、ソースがそれを生み出さなかったかのように扱われます。対応するキーは置き換えも削除もせず、フローはそのキーの最新の有効な行に戻されます。 このフローは、正の total_amountでない更新をドロップします。

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")