進階自動疾病控制中心主題

除了基本 AUTO CDC 功能和 AUTO CDC FROM SNAPSHOT API,你還可以在目標資料表上執行 DML、讀取 CDC 目標的變更資料串流、監控處理指標、套用部分更新,並透過雙時式儲存追蹤變更。 API 的介紹請參見 AUTO CDC ,「AUTO CDC API 範例:使用管線簡化變更資料擷取」。

新增、變更或刪除目標串流資料表中的資料

如果您的管線將資料表發佈至 Unity 目錄,您可以使用 資料操作語言 (DML) 陳述式,包括插入、更新、刪除和合併陳述式,來修改陳述式所 AUTO CDC ... INTO 建立的目標串流資料表。

備註

  • 不支援修改串流數據表之數據表架構的 DML 語句。 請確定您的 DML 語句不會嘗試演進數據表架構。
  • 更新串流數據表的 DML 語句只能在共用的 Unity 目錄叢集或 SQL 倉儲中使用 Databricks Runtime 13.3 LTS 和更新版本來執行。
  • 因為串流處理需要只能附加的資料來源,如果您的處理需要從具有變更的來源串流資料表進行串流處理(例如 DML 語句),請在讀取來源串流資料表時設定 skipChangeCommits 旗標。 設定 skipChangeCommits 時,會忽略刪除或修改源數據表上記錄的交易。 如果您的處理不需要串流數據表,您可以使用具象化視圖(不受僅附加限制)作為目標數據表。

由於管線會使用指定的 SEQUENCE BY 欄位,並將適當的排序值傳播到目標資料表的 __START_AT 和 __END_AT 欄位(針對 SCD Type 2),因此您必須確保 DML 陳述式對這些欄位使用有效的值,以維持記錄的正確排序。 請參閱 AUTO CDC 的運作方式。

如需搭配串流資料表使用 DML 陳述式的詳細資訊,請參閱 新增、變更或刪除串流資料表中的資料。

下列範例會插入起始序列為 5 的有效記錄:

INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);

小提示

如果你需要在 SCD Type 2 目標表中重新命名 __START_AT 和 __END_AT 欄位(例如,為了符合下游結構需求),請建立一個針對目標資料表的檢視:

CREATE VIEW my_employees_view AS
SELECT
  *,
  __START_AT AS valid_from,
  __END_AT AS valid_to
FROM my_scd2_target_table;

從 AUTO CDC 目標資料表讀取變更資料摘要

在 Databricks Runtime 15.2 和更新版本中,您可以從作為 AUTO CDC 或 AUTO CDC FROM SNAPSHOT 查詢目標的串流資料表中,讀取變更資料提要,其方式與從其他 Delta 表讀取變更資料提要的方法相同。 從目標串流資料表讀取變更資料提要需要下列項目:

  • 目標串流資料表必須發佈至 Unity 目錄。 請參閱 使用 Unity 目錄搭配管線。
  • 若要從目標串流數據表讀取變更數據摘要,您必須使用 Databricks Runtime 15.2 或更新版本。 若要讀取不同管線中的變更數據摘要,管線必須設定為使用 Databricks Runtime 15.2 或更新版本。

您可以用與讀取其他 Delta 資料表的變更資料饋送相同的方式,讀取在 Lakeflow 管線中建立的目標串流資料表的變更資料饋送。 若要進一步了解如何使用 Delta 變更資料摘要功能(包括 Python 和 SQL 範例),請參閱 在 Azure Databricks 上使用變更資料摘要。

備註

變更資料摘要記錄包含識別變更事件類型的 中繼資料 。 在表格中更新記錄時,相關聯變更記錄的中繼資料通常包括 _change_type 設定為 update_preimage 和 update_postimage 事件的值。

不過, _change_type 如果對目標串流資料表進行更新,包括變更主索引鍵值,則值會有所不同。 當變更包含主索引鍵的更新時, _change_type 中繼資料欄位會設定為 insert 和 delete 事件。 當使用 UPDATE 或 MERGE 陳述式手動更新其中一個鍵欄位時,或者對於 SCD 類型 2 表格,當 __start_at 欄位變更以反映較早的啟動順序值時,可能會發生主鍵的變更。

查詢會決定主鍵值,這些值在 SCD 類型 1 和 SCD 類型 2 的處理過程中有所不同:

SCD 類型 主要金鑰
SCD 類型 1,以及管線 Python 介面 主鍵是函式中keys參數create_auto_cdc_flow()的值。 對於 SQL 介面,主鍵是由KEYS子句在AUTO CDC ... INTO陳述式中定義的欄。
SCD 類型 2 主鍵為 keys 參數或 KEYS 子句加上操作的 coalesce(__START_AT, __END_AT) 回傳值,其中 __START_AT 和 __END_AT 分別是目標串流表中的對應欄位。 此方法在可用時使用,且__START_AT當__END_AT為空(例如初始記錄)時使用__START_AT。

從實體化視圖讀取變更資料流

Important

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

你可以從 Lakeflow 管線或 Databricks SQL 中建立的實體化視圖讀取變更資料饋送。 利用此資料將實體化檢視變更複製到Azure Databricks以外的目的地,或保留實體化檢視變更的歷史以供稽核與報告。

實體化檢視使用自動變更資料串流,因此你不會直接開啟變更資料串流。 相反地,你只需滿足以下要求,就能在每個具體化的視圖上啟用變更資料串流。 請參見 自動變更資料串流。

  • 要讀取變更資料饋送,您必須使用 Databricks Runtime 18 LTS 或以上版本,適用於經典運算、無伺服器運算或 Databricks SQL。

  • 具體化視圖、產生它的管線或讀取它的管線必須使用該 PREVIEW 通道。

  • 實體化視圖必須啟用列追蹤功能。 無伺服器計算的實體化視圖預設啟用了列追蹤功能。 請參見 Azure Databricks 中的列追蹤。 要檢查實體化視圖是否啟用列追蹤,執行:

    SHOW TBLPROPERTIES my_mv ('delta.enableRowTracking');
    
  • 若要從實體化檢視讀取異動資料摘要,請在管線或實體化檢視上啟用外部中繼資料標記。 有關說明,請參見 「如何啟用資料集存取」。

你可以使用 table_changes() 函式、串流讀取或 readChangeFeed 選項,以與其他 Delta 資料表相同的方式讀取具體化檢視表的變更資料摘要。 關於 SQL 和 Python 的語法與範例,請參見 Use change data feed on Azure Databricks。

你可以在 Databricks SQL 具體化檢視或串流資料表內讀取具體化檢視的變更資料摘要:

CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT * FROM STREAM my_mv WITH (readChangeFeed=true)

Limitations

除了 自動變更資料饋送限制之外,當您從實體化檢視讀取變更資料饋送時,還適用下列限制:

  • 當實體化檢視完全重寫時,變更資料串流仍包含未變更的資料列,且不會將同一列的多個更新合併為單一事件。 若要將這些資料篩除,請依所有欄位分組來彙整變更資料摘要,以找出具有相同資料列值的插入與刪除。
  • 只有 Azure Databricks 能查詢變更資料流以取得實體化的視圖。 外部 Delta Lake 和 Iceberg 用戶端則無法。
  • 在 Lakeflow 管線中,你只能從另一個不同的管線讀取具體化檢視的變更資料饋送,而且該管線必須使用 PREVIEW 通道。 不支援在建立該實體化檢視的同一個管線中讀取其變更資料摘要。
  • 你無法從實體化的視圖建立向量搜尋索引。

取得管線中 CDC 查詢所處理之記錄的相關資料

備註

下列指標只會由 AUTO CDC 查詢擷取,而不是由 AUTO CDC FROM SNAPSHOT 查詢擷取。

查詢 AUTO CDC 會收集以下指標:

  • num_upserted_rows:更新期間更新插入資料集的輸出列數。
  • num_deleted_rows:更新期間從資料集中刪除的現有輸出列數。

num_output_rows 指標(非 CDC 流程的輸出)不會針對 AUTO CDC 查詢擷取。

套用部分更新

當來源只傳送變更的欄位時,AUTO CDC必須區分變更記錄中缺失的欄位(應保持目標值不變)與明確設定為 null的欄位,該欄位應覆蓋目標值。null 預設情況下, IGNORE NULL UPDATES 將 每個 null 視為「不更新」標記,因此無法套用明確 null的 。 為了解決這個歧義,請選擇以下三種方法之一:

方法 何時使用 行為
IGNORE NULL UPDATES ON columnList 一組小型固定欄位應忽略 null 值,而其他欄位則會套用明確 null 值。 當輸入值為 null時,列出的欄位仍保留其現有目標值。 其他欄位則會套用明確 null 的值。
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList) 大多數欄位應該忽略 null 這些值,只有少數欄位會明確套用 null 值。 所列欄位套用明確的 null 值。 當輸入值為 null時,其他欄位皆保留其現有目標值。
COLUMNS TO UPDATE 每個變更記錄會更新不同的欄位集合,或是可更新欄位的集合會隨時間改變。 來源欄位會指定每個變更記錄要更新的欄位。 列出的欄位皆從來源寫入,包含明確的 null 值。 未列出的欄位則保留其現有的目標值。

COLUMNS TO UPDATE 無法與 IGNORE NULL UPDATES合併,且不支援雙時空表。

經驗法則是,當產生者知道每筆記錄中哪些欄位已變更,並可在來源欄位中攜帶該資訊時,請選擇 COLUMNS TO UPDATE;例如,當多個產生者寫入同一個來源,或可更新的欄位集合會隨時間擴大時。 選擇 IGNORE NULL UPDATES ON 當管線擁有者事先知道固定可更新欄位集合,並偏好以管線程式碼控制它們時。

以下範例使用一個名為 columnsToUpdate 的來源欄位來控制變更記錄更新的欄位,包括明確設定為 null的欄位:

Python

from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
  target = "target",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "sequenceNum",
  stored_as_scd_type = 1,
  columns_to_update = "columnsToUpdate"
)

SQL

CREATE OR REFRESH STREAMING TABLE target;

CREATE FLOW apply_cdc AS AUTO CDC INTO
  target
FROM
  stream(cdc_source)
KEYS
  (id)
SEQUENCE BY
  sequenceNum
STORED AS
  SCD TYPE 1
COLUMNS TO UPDATE
  columnsToUpdate;

如需完整的參數參考資料,請參閱 AUTO CDC INTO (pipelines) 及 create_auto_cdc_flow。

雙時態自動 CDC

Important

Bitemporal AUTO CDC 目前處於 Beta 階段。

SCD 第 1 型和第 2 型屬於單時態:它們沿著單一時間維度追蹤變更。 Bitemporal 擴展了 SCD 第二型歷史,追蹤兩個時間維度的變化,並區分兩種視角:

  • 商務時間:事件實際發生的時間。
  • 系統時間:系統記錄或接收事件的時間。

與 SCD Type 2 類似,雙時態會保留記錄的完整歷史。 它增加了第二條時間軸,讓你能重建數據顯示的內容以及系統過去任何時候的信念。

例如,對沖基金會從來源系統接收股票數據。 Acme Corp 的股價將於 1 月 1 日變動,但基金直到 1 月 5 日才會接收該更新。 Bitemporal AUTO CDC 讓基金能回答兩個明確的問題:Acme Corp 在 1 月 1 日(營業時間)的實際股價,以及系統在 1 月 3 日基金做出交易決策時所相信的價格(系統時間)。 能夠區分這些時間軸,對於審計、監管報告及財務決策非常有幫助。

若要啟用雙時間處理,請設定 STORED AS BITEMPORAL(SQL)或 stored_as_scd_type="bitemporal"(Python),並將 SEQUENCE BY 用於業務時間欄位,將 SYSTEM SEQUENCE BY 用於系統時間欄位。 目標資料表會新增 __SYSTEM_START_AT 和 __SYSTEM_END_AT 欄位,以及 SCD 第 2 類型的 __START_AT 和 __END_AT 欄位。 語法細節請參見 AUTO CDC INTO (pipelines) 或 create_auto_cdc_flow。

雙時態 AUTO CDC 範例

以下範例從一小組合成的 CDC 事件中建立雙時態目標表。 bt欄位顯示業務時間,欄位st顯示系統時間。

Python

from pyspark import pipelines as dp

# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")

@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
  return spark.createDataFrame(
    [
      (1, "x10", "y10", 10, 100),
      (1, "x20", "y20", 20, 200)
    ],
    schema="id INT, x STRING, y STRING, bt INT, st INT",
  )

# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")

dp.create_auto_cdc_flow(
  target = "target_bitemporal",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "bt",
  system_sequence_by = "st",
  stored_as_scd_type = "bitemporal"
)

SQL

-- Source: synthetic CDC events
CREATE OR REFRESH STREAMING TABLE cdc_source_sql;

CREATE FLOW cdc_source_sql AS INSERT INTO ONCE
  cdc_source_sql BY NAME
SELECT * FROM VALUES
  (1, 'x10', 'y10', 10, 100),
  (1, 'x20', 'y20', 20, 200)
  AS t(id, x, y, bt, st);

-- Target: bitemporal table
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;

CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
  target_bitemporal_sql
FROM
  stream(cdc_source_sql)
KEYS
  (id)
SEQUENCE BY
  bt
SYSTEM SEQUENCE BY
  st
STORED AS
  BITEMPORAL;

以下變更序列展示了雙時空表格如何記錄單一公司的插入、更新、錯序更新與刪除。 序列欄位產生 __START_AT 與 __END_AT (業務時間)欄位,系統序列欄位產生 __SYSTEM_START_AT 與 __SYSTEM_END_AT (系統時間)欄位:

Column Description
__START_AT 這場爭執生效的營業時間。
__END_AT 此列有效性結束的業務時間。 null 如果是無限期有效的話。
__SYSTEM_START_AT 此列的資料及其業務時間區間已知屬實時的系統時間點。
__SYSTEM_END_AT 可確定此列資料及其業務時間區間已失效的系統時間。 null 如果被無限期地證明為真。

系統能處理兩條時間線上任意順序發生的事件。 當事件以比已處理事件更早的業務時間或系統時間抵達時,系統會修正受影響的歷史,而非僅附加到末尾。

變更 1:插入

公司A於2025年7月18日10:01:00(營業時間)加入,但直到系統時間10:05:00才被納入。

輸入:

公司識別碼 數據點 排序 系統排序 運算
A XFv1 7/18/2025 10:01:00 7/18/2025 10:05:00 INSERT

輸出:

公司識別碼 數據點 __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 零 7/18/2025 10:05:00 零

XFv1 從 10:01:00 開始有效,且無已知結束時間。 系統在系統時間 10:05:00 得知這項事實,且其結束時間尚未知。

變更 2:更新

A 公司於 2025年7月18日 12:15:43(業務時間)更新,系統於 12:20:00(系統時間)處理該事件。 系統會同時保留在得知更新之前原先認定的內容,以及在匯入更新後更正過的業務歷史。

輸入:

公司識別碼 數據點 排序 系統排序 運算
A XFv2 7/18/2025 12:15:43 7/18/2025 12:20:00 UPDATE

輸出:

公司識別碼 數據點 __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 零 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 零
A XFv2 7/18/2025 12:15:43 零 7/18/2025 12:20:00 零

XFv1 被認為自 10:01:00 起有效且無明確結束時間,系統也維持此信念至 10:05:00 至 12:20:00。 XFv1 現已確認僅有效至 12:15:43,而修正後的歷史記錄自系統時間 12:20:00 起生效,目前已知無結束時間。 XFv2 從 12:15:43 開始有效,無已知結束時間,並於系統時間 12:20:00 學習。

變更 3:亂序更新

出現一個故障更新,顯示公司 A 實際上於 2025/7/18 12:05:00(營業時間)更新,但直到系統時間 12:25:00 才被接收。 當更新在系統時間較晚但業務時間之前時,系統會修正歷史業務時間,並保留錯誤更新前的判斷與修正後的歷史。

輸入:

公司識別碼 數據點 排序 系統排序 運算
A XFv3 7/18/2025 12:05:00 7/18/2025 12:25:00 UPDATE

輸出:

公司識別碼 數據點 __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 零 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 零
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 零
A XFv2 7/18/2025 12:15:43 零 7/18/2025 12:20:00 零

XFv1 被認為有效於 10:01:00 至 12:15:43,而該信念現在在系統時間中也有效,直到 12:25:00。 新更新將 XFv1 的業務有效期修正為 12:05:00,該修正歷史自系統時間 12:25:00 起生效。 XFv3 現已知在 12:05:00 至 12:15:43 期間有效,而此一判定在系統時間中自 12:25:00 起有效,且目前沒有已知的結束時間。

變更四:刪除

A 公司於 2025 年 7 月 18 日 12:30:00 被刪除,而系統於 12:30:00 處理該事件。 由於刪除操作代表實體業務存在的終結,系統不會產生替代列。 XFv2 以兩列形式出現,完整記錄公司解散及系統得知刪除時的審計軌跡。

輸入:

公司識別碼 數據點 排序 系統排序 運算
A XFv2 7/18/2025 12:30:00 7/18/2025 12:30:00 DELETE

輸出:

公司識別碼 數據點 __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 零 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 零
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 零
A XFv2 7/18/2025 12:15:43 零 7/18/2025 12:20:00 7/18/2025 12:30:00
A XFv2 7/18/2025 12:15:43 7/18/2025 12:30:00 7/18/2025 12:30:00 零

XFv2 從 12:15:43 有效且無已知結束時間,系統則在 12:20:00 至 12:30:00 期間保持此信念。 刪除作業被匯入後,可確定 XFv2 僅在 12:30:00 之前有效,而修正後的歷史記錄自系統時間 12:30:00 起生效。

哪些資料物件用於管線中的 CDC 處理?

當您在 Hive 中繼存放區中宣告目標資料表時,會建立兩個資料結構:

  • 使用目標表格所指派名稱的檢視。
  • 管線用來管理 CDC 處理的內部備份資料表。 此表格的命名方式是將 __apply_changes_storage_ 添加在目標表格名稱之前。

例如,如果你宣告了一個名稱為 dp_cdc_target 的目標資料表,你會在元儲存庫中看到一個名為 dp_cdc_target 的檢視表和一個名為 __apply_changes_storage_dp_cdc_target 的資料表。 查詢視圖以存取處理過的資料。 不要直接修改支援資料表。

備註

這些資料結構僅適用於 AUTO CDC 處理,不適用於 AUTO CDC FROM SNAPSHOT 處理。 它們僅適用於 Hive 元資料庫,不適用於 Unity 資料目錄。