使用 REPLACE USING 流程進行部分快照替換

Important

這項功能位於 測試版 (Beta) 中。

REPLACE USING FLOW 會讓目標資料表與串流來源保持同步:它會替換所有與指定鍵欄位相符的資料列,並保持其他資料不變。

欄位 SEQUENCE BY 會排序更新,確保即使更新順序錯亂,結果依然正確。 對於每個鍵,最高序列勝出,且較低序列的列永遠不會覆蓋目標中已存在的較高序列。 使用相同鍵與相同序列的列會被附加,而非替換。

REPLACE USING 的運作方式

考慮一個事件表,記錄兩個地區的點擊與轉換事件,順序如下 seq

region_id device_type 事件類型 序列
1 iOS click 1
1 Android 轉換 1
2 iOS click 1
2 桌面 click 1

一個 REPLACE USING (region_id) SEQUENCE BY seq 流程會接收區域 1 和 3 的這些更新。 第二區沒有更新:

region_id device_type event_type 序列
1 iOS click 2
1 Android 轉換 2
1 桌面 click 2
3 iOS click 1
3 桌面 click 2

目標為:

region_id device_type event_type 序列 結果
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 列不會被加入,因為只套用了金鑰最高序列。

要求

使用 REPLACE 的流程須符合以下要求:

  • REPLACE USING 流程可在 Databricks Runtime 18.2 及以上版本上執行,並可使用傳統或無伺服器運算。 Databricks 推薦 Unity Catalog。
  • 來源必須是串流來源。 REPLACE USING 會拒絕非串流來源。
  • 你必須指定至少一個關鍵欄位,且必須精確指定一 SEQUENCE BY 欄。

何時使用 REPLACE USING 流程

Lakeflow 管道提供三種會覆寫現有資料列的流程。 根據您的來源資料樣貌,以及其如何辨識要取代的資料列來進行選擇:

  • 當您的來源是一連串以欄位作為鍵的部分快照時,請使用 REPLACE USING。 REPLACE USING 只會覆寫在傳入資料中有相符項目的資料,其他所有資料則保持不變。 它不需要主金鑰。
  • 當您的來源是帶有明確插入更新刪除操作的變更資料擷取(CDC)資料流,或您需要緩慢變更維度(SCD)類型 2 歷史時,請使用 AUTO CDC。 AUTO CDC 也需要真正的主金鑰。 請參閱 AUTOTO CDC API:使用管線簡化變更資料擷取
  • 來源是快照,而你想重新計算並覆寫以述詞選取的目標資料表範圍(例如最近 7 天的資料),並將此作業作為批次操作執行時,請使用 REPLACE WHERE 它不需要主金鑰。 請參見 使用 REPLACE WHERE 流程的批次處理

建立一個替換使用流程

請用 SQL 或 Python 定義替換使用流程。

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 欄位可讓結果不受更新到達順序的影響。 只有當某一列的序列大於該鍵目前已儲存的序列時,這一列才會套用到該鍵;因此,若延遲到達或重放的列比目前值還舊,則會被忽略。 更新中沒有的金鑰則保持不變。

遵循以下做法,讓替換行為可預測:

練習 原因
使用嚴格依照金鑰版本增加的序列,例如時間戳、版本號或日誌偏移量。 具有相同鍵和相同序列的兩列都會被保留,因此該鍵會出現重複列。
使用非零序列。 空序列可能導致行為未定義。

Expectations

「REPLACE」使用流程符合預期。 warnfail 的行為與其他流程中的情況相同:warn 會保留違規的資料列並記錄該違規,而 fail 則會停止更新。 請參閱 使用管線期望來管理資料品質

drop 期望會將違規的資料列視為如同來源端從未產生過它。 刪除的列不會替換、刪除或修改目標表中的相符鍵:

  • 丟棄發生在去重之前,因此流程會保留該金鑰的最新有效版本。
  • 如果每個輸入的密鑰列都被刪除,該密鑰現有的列則保持不變。
  • 由於刪除的列不會設定序列底線,即使後續的更新序列低於刪除列的,仍會生效。

Limitations

替換使用流程有以下限制:

  • REPLACE USING 僅支援每個目標資料表的單一流程。 在同一目標上,不支援將 REPLACE USING 與其他流程類型合併。
  • 目標資料表必須在管線內建立。
  • 來源必須是串流來源。
  • 你必須指定至少一個鍵欄位和一個 SEQUENCE BY 欄位。 鍵欄位不能重複,且每個鍵欄位的類型必須是可排序的。 原子型態,如整數、字串和日期,可以是鍵,而 MAPVARIANT 則不能。
  • 關於獨立串流資料表,請參見以 REPLACE USING 流程進行部分快照替換 ,以了解語法差異。

Examples

以下範例是從 samples.wanderbricks.booking_updates 讀取,這是一個在每個啟用 Unity Catalog 的工作區中都可使用的預約狀態變更範例資料表。 每筆預訂在每次變更時都會出現一次,因此 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")