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」使用流程符合預期。
warn 和 fail 的行為與其他流程中的情況相同:warn 會保留違規的資料列並記錄該違規,而 fail 則會停止更新。 請參閱 使用管線期望來管理資料品質。
drop 期望會將違規的資料列視為如同來源端從未產生過它。 刪除的列不會替換、刪除或修改目標表中的相符鍵:
- 丟棄發生在去重之前,因此流程會保留該金鑰的最新有效版本。
- 如果每個輸入的密鑰列都被刪除,該密鑰現有的列則保持不變。
- 由於刪除的列不會設定序列底線,即使後續的更新序列低於刪除列的,仍會生效。
Limitations
替換使用流程有以下限制:
- REPLACE USING 僅支援每個目標資料表的單一流程。 在同一目標上,不支援將 REPLACE USING 與其他流程類型合併。
- 目標資料表必須在管線內建立。
- 來源必須是串流來源。
- 你必須指定至少一個鍵欄位和一個
SEQUENCE BY欄位。 鍵欄位不能重複,且每個鍵欄位的類型必須是可排序的。 原子型態,如整數、字串和日期,可以是鍵,而MAP和VARIANT則不能。 - 關於獨立串流資料表,請參見以 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")