變更資料摘要(CDF)會追蹤 Delta Lake 資料表或 Apache Iceberg v3 資料表不同版本之間的列層級變更。 每個變更記錄包含列資料以及指示該列是否入、更新或刪除的元資料。
你可以將變更資料饋送用於常見的資料使用情境,包括:
- 僅處理自上次管線執行以來已變更資料列的增量 ETL 管線。
- 用於追蹤資料修改以符合合規與治理要求的稽核軌跡。
- 資料複製工作負載,會將變更同步到下游資料表、快取或外部系統。
Azure Databricks 支持兩種方法:
- 自動變更資料導向:在資料表讀取期間,利用列系譜元資料計算變更。 這不需要個別資料表設定,且適用於 Delta Lake 和 Apache Iceberg v3 資料表。 請參見 自動變更資料串流。
- 舊有變更資料流:在資料表寫入時實現變更。 只支援 Delta Lake 表格。 需要個別的表格配置。 請參閱 Delta Lake 的 Legacy 變更資料饋送。
自動變更資料串流
自動變更資料摘要會在查詢時而非在寫入時計算資料列層級的變更,方式是使用 Delta Lake 資料表上的資料列追蹤,以及 Apache Iceberg v3 資料表上的資料列譜系。 與傳統變更資料串流不同,它不需要個別資料表的設定。 任何符合需求的表格都會自動支援。 請參閱 需求。
由於在每次執行 MERGE INTO 和 UPDATE 寫入作業時不會計算變更,因此相較於舊版變更資料摘要,自動變更資料摘要可提升寫入效能並降低儲存成本。
自動變更資料串流使用與舊有變更資料串流相同的 table_changes()readChangeFeed API,並支援批次查詢、結構化串流及 Databricks-to-Databricks OpenSharing。 請參閱 在批次查詢中讀取變更及 以增量方式處理變更資料。
需求規格
- Databricks 執行時間 19 或以上
- 一種已註冊於 Unity 目錄的支援表格格式:
- 一個可管理的表格,採用 Delta Lake 格式,啟用行追蹤功能,或 Iceberg v3 格式。
- 一個以 Delta Lake 格式呈現並啟用行追蹤的外部表格。
註記
變更資料串流不包含在 Apache Iceberg 規範中。Azure Databricks 讀取器可以查詢 Apache Iceberg v3 資料表的自動變更資料串流,但外部 Iceberg 讀取器則不行。 請參閱 冰山數據表規格。
對於 Delta Lake,只有 Azure Databricks 讀取器可以查詢自動變更資料串流。
使用變更資料摘要
要使用變更資料串流,請確認你使用的表格符合要求。 請參閱 需求。
若要以批次方式讀取變更資料摘要,請執行下列操作:
Python
spark.read \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.table("<table_name>")
Scala
spark.read
.option("readChangeFeed", "true")
.option("startingVersion", 0)
.table("<table_name>")
SQL
SELECT * FROM table_changes('<table_name>', 0)
如需進一步了解如何以批次方式讀取變更資料導送,請參閱「在批次查詢中讀取變更」。
若要以串流方式讀取變更資料饋送,請執行下列步驟:
Python
(spark.readStream
.option("readChangeFeed", "true")
.table("<table_name>")
)
Scala
spark.readStream
.option("readChangeFeed", "true")
.table("<table_name>")
如需進一步了解變更資料摘要的串流讀取,請參閱 「以增量方式處理變更資料」。
從舊有變更資料串流遷移
要將 Delta Lake 資料表從舊有變更資料串流遷移到自動變更資料饋送,請執行以下步驟:
- 確認你的桌子符合 要求。
- 請執行以下指令關閉舊有變更資料串流:
ALTER TABLE <table_name> UNSET TBLPROPERTIES ('delta.enableChangeDataFeed');
你不能同時使用舊有資料和自動變更資料串流。
變更資料饋源架構
當你從資料表的變更資料串流讀取時,查詢會使用最新版本的資料表架構。 Azure Databricks 支援大部分的結構變更與演化操作,但帶有欄位映射的資料表有其限制。 參見 帶有欄位映射的表格。
除了來自 Delta Lake 資料表結構的資料欄位外,變更資料串流還包含用以識別變更事件類型的元資料欄位:
| 欄位名稱 | 類型 | 價值 |
|---|---|---|
_change_type |
String | 包含:insert、、update_preimageupdate_postimage、delete。preimage 是更新前的值, postimage 是更新後的值。 |
_commit_version |
Long | 包含:包含該變更的 Delta 記錄或資料表版本。 |
_commit_timestamp |
時間戳記 | 包含:與提交建立時間相關聯的時間戳記。 |
如果結構中包含與這些元資料欄位名稱相同的欄位,你就無法在資料表中使用變更資料串流。 在開啟變更資料串流前,先重新命名資料表中的欄位以解決此衝突。
以增量方式處理變更資料
Databricks 建議你結合 Structured Streaming 使用變更資料流,以增量式處理資料表的變更。 您必須使用 Azure Databricks 的結構化串流來自動追蹤數據表變更數據摘要的版本。 關於使用 SCD 類型 1 或類型 2 表格進行 CDC 處理,請參見 The AUTO CDC APIS:簡化使用管線變更資料擷取。
當串流首次啟動時,變更資料饋送會先以 INSERT 記錄的形式傳回資料表的最新快照,接著再將未來的變更作為變更資料傳回。 變更資料導流會同時將變更資料與新資料列提交到資料表交易日誌。
若要設定串流以讀取資料表的變更資料摘要,請將選項 readChangeFeed 設為 true,如下:
Python
(spark.readStream
.option("readChangeFeed", "true")
.table("myTable")
)
Scala
spark.readStream
.option("readChangeFeed", "true")
.table("myTable")
速率限制
Azure Databricks 在讀取變更資料時支援速率限制(maxFilesPerTrigger、maxBytesPerTrigger)和 excludeRegex。 欲了解完整的串流三角洲湖選項列表,請參見 三角洲湖。
你也可以選擇指定起始版本,詳見 「指定起始版本」。 對於起始快照以外的版本,速率限制會以原子方式套用於整個提交。 目前批次要麼包含整個提交,要麼將該提交延至下一個批次。
重賽積分榜歷史
變更資料串流並非用來作為表格所有變更的永久記錄。 它只記錄在啟用變更資料串流後發生的變更。 你可以開始新的串流讀取,以擷取目前版本及之後的所有變更。
變更資料饋送中的紀錄是暫時性的,且僅在指定的保留期間內可存取。 交易日誌會定期移除資料表版本及其對應的變更資料饋送版本。 當某個版本被移除時,你就無法再讀取該版本的變更資料饋送。
封存變更資料,以永久保留歷史記錄
如果你的使用情境需要永久保存資料表所有變更的歷史,可以用增量邏輯將變更資料流的紀錄寫入新資料表。
以下範例示範如何使用 trigger.AvailableNow 以批次工作負載的方式處理可用資料,用於稽核或完整變更重播:
Python
(spark.readStream
.option("readChangeFeed", "true")
.table("source_table")
.writeStream
.option("checkpointLocation", <checkpoint-path>)
.trigger(availableNow=True)
.toTable("target_table")
)
Scala
spark.readStream
.option("readChangeFeed", "true")
.table("source_table")
.writeStream
.option("checkpointLocation", <checkpoint-path>)
.trigger(Trigger.AvailableNow)
.toTable("target_table")
指定起始版本
要讀取特定點的變更,請使用時間戳記或版本號指定起始版本。 批次讀取需要起始版本。 你也可以指定結束版本,以限制範圍。 若要進一步了解資料表歷史,請參見 時間旅行。
當你配置使用變更資料流的結構化串流工作負載時,指定起始版本可能會影響處理效能:
- 新的資料處理管線通常會受益於預設行為,即在資料流開始時將資料表中所有現有紀錄記錄為
INSERT操作。 - 如果您的目標數據表已包含所有具有適當變更的記錄,請指定起始版本以避免將源數據表狀態當做
INSERT事件處理。
以下範例展示了如何在檢查點損壞的情況下從串流失敗中恢復。 在此範例中,假設有下列條件:
- 在數據表建立時,源數據表上已啟用變更數據摘要。
- 目標下游資料表處理所有變更,直到版本 75 為止。
- 源數據表的版本歷程記錄適用於版本70和更新版本。
當你定義寫入串流到現有目標資料表時,必須指定一個新的檢查點位置:
Python
(spark.readStream
.option("readChangeFeed", "true")
.option("startingVersion", 76)
.table("source_table")
.writeStream
.option("checkpointLocation", "<new-checkpoint-path>")
.toTable("target_table")
)
Scala
spark.readStream
.option("readChangeFeed", "true")
.option("startingVersion", 76)
.table("source_table")
.writeStream
.option("checkpointLocation", "<new-checkpoint-path>")
.toTable("target_table")
這很重要
如果你指定了起始版本,但該版本在資料表歷程記錄中不存在,串流將無法從新的檢查點啟動。 由於受管理資料表會自動清理歷史版本,所有指定的起始版本最終都會被刪除。
請參閱 重播桌歷史。
讀取批次查詢中的變更
你可以使用批次查詢語法,從特定版本開始讀取所有變更,或在指定版本範圍內讀取變更,如下:
- 將版本指定為整數,時間戳記為格式中的
yyyy-MM-dd[ HH:mm:ss[.SSS]]字串。 - 開始和結尾版本皆包含所有內容。 要從起始版本讀取到最新版本,請只指定起始版本。
- 如果你指定的版本早於變更資料摘要啟用的時間,就會引發錯誤。
要使用具有起始與結束版本選項的批次讀取,請執行以下步驟:
SQL
要從版本 0 閱讀到 10,請執行以下操作:
SELECT * FROM table_changes('tableName', 0, 10)
若要讀取兩個時間戳記版本之間的資料,請執行下列步驟:
--
SELECT * FROM table_changes('tableName', '2021-04-21 05:45:46', '2021-05-21 12:00:00')
要從起始版本讀到最新版本,請做以下操作:
SELECT * FROM table_changes('tableName', 0)
要讀取名稱中帶有特殊字元的資料表的變更,請執行以下操作:
SELECT * FROM table_changes('`schema`.`dotted.tableName`', '2021-04-21 06:45:46', '2021-05-21 12:00:00')
請參閱 table_changes 資料表值函式。
Python
要從版本 0 閱讀到 10,請執行以下操作:
spark.read \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.option("endingVersion", 10) \
.table("myDeltaTable")
要在兩個時間戳之間讀取,請執行以下操作:
spark.read \
.option("readChangeFeed", "true") \
.option("startingTimestamp", '2021-04-21 05:45:46') \
.option("endingTimestamp", '2021-05-21 12:00:00') \
.table("myDeltaTable")
要從起始版本讀到最新版本,請做以下操作:
spark.read \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.table("myDeltaTable")
Scala
要從版本 0 閱讀到 10,請執行以下操作:
spark.read
.option("readChangeFeed", "true")
.option("startingVersion", 0)
.option("endingVersion", 10)
.table("myDeltaTable")
要在兩個時間戳之間讀取,請執行以下操作:
spark.read
.option("readChangeFeed", "true")
.option("startingTimestamp", "2021-04-21 05:45:46")
.option("endingTimestamp", "2021-05-21 12:00:00")
.table("myDeltaTable")
要從起始版本讀到最新版本,請做以下操作:
spark.read
.option("readChangeFeed", "true")
.option("startingVersion", 0)
.table("myDeltaTable")
處理超出範圍的版本
預設情況下,如果你指定一個版本或時間戳超過上次提交,查詢會回傳錯誤 timestampGreaterThanLatestCommit。
在 Databricks Runtime 11.3 LTS 及以上版本中,您可以啟用對超出範圍版本的容忍度,方法如下:
SET spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = true;
啟用此配置時,查詢會回傳不同的結果如下:
- 起始版本或超過最後一次提交的時間戳記會回傳空結果。
- 結束版本或晚於最後一次提交的時間戳記,會回傳從起始點到最後一次提交之間的所有變更。
三角湖的舊有變更資料饋送
舊版變更資料饋送需要為個別的 Delta Lake 資料表逐一手動設定。 由於變更資料流未包含在 Apache Iceberg 規範中,因此不支援 Apache Iceberg 資料表。 Databricks 建議你遷移到自動變更資料串流。 請參見「從舊有變更資料串流遷移」。
當舊有變更資料導向開啟時,執行時會記錄所有寫入資料表的資料 變更事件 。 其中包含列資料,以及指出指定資料列是否已插入、刪除或更新的元資料。
舊版變更資料饋送使用與自動變更資料饋送相同的 readChangeFeed 和 table_changes() 讀取 API。 請參見「 增量處理變更資料 」及「 在批次查詢中讀取變更」。
開啟舊有變更資料串流
你必須在各個資料表上明確開啟舊有變更資料串流。 請使用下列其中一個方法:
新增資料表
請在delta.enableChangeDataFeed = true指令中設定CREATE TABLE的 table 屬性。
CREATE TABLE student (id INT, name STRING, age INT)
TBLPROPERTIES (delta.enableChangeDataFeed = true)
註記
如果你在某段時間內關閉舊有變更資料串流,然後再重新開啟,該期間將無法查詢。 使用自動變更資料流在區間內查詢變更。 請參見 自動變更資料串流。
現有的資料表
請在delta.enableChangeDataFeed = true指令中設定ALTER TABLE的 table 屬性。
ALTER TABLE myDeltaTable
SET TBLPROPERTIES (delta.enableChangeDataFeed = true)
儲存考量
受管理資料表能有效記錄資料變更,並可能使用其他功能來優化儲存配置。
使用舊有變更資料流時,您必須考慮以下儲存行為:
- 你可能會看到儲存成本略有增加,因為變更可能會被記錄在獨立檔案中。
- 某些操作,如僅插入或完整分割區刪除,不會產生變更資料檔。 Azure Databricks 直接從交易日誌計算變更資料流。
- 變更資料檔案會使用資料表的保留政策。
VACUUM命令會刪除變更資料檔案,而來自交易記錄的變更則會使用檢查點保留原則。
Databricks 建議你不要直接查詢變更資料檔來重建變更資料串流。 一定要使用 Delta Lake 和 Apache Iceberg API。
局限性
請考慮變更資料串流的以下限制:
具有欄位映射的表格
在 Delta Lake 表格啟用欄位映射後,你可以刪除或重新命名欄位,而無需重寫資料檔案。 請參閱 使用 Delta Lake 欄位映射重新命名和刪除欄位。
然而,變更資料串流在非加法結構變更後仍有限制。 非加法結構變更包括以下操作:
- 重新命名或刪除欄位。
- 更改欄位資料型別。
- 改變欄位的空值,例如。
ALTER COLUMN ... SET NOT NULL參見NOT NULL限制。
你無法讀取在發生非加性結構描述變更之交易或範圍中的變更資料流。
為了允許在指定批次讀取範圍前後進行非加法式結構變更,查詢使用該範圍最終版本的結構,而非最新的表格版本。 即使版本範圍跨越非加法結構變更,查詢仍會失敗。