Note
Lakebase 變更資料饋送功能目前處於 公開預覽階段。
什麼是 Lakebase 變更資料串流?
Lakebase 引入原生變更資料饋送(CDF),解鎖您的營運資料,用於下游的管線、模型與應用程式。 Lakebase Postgres 表格上的每次插入、更新和刪除都會從預先寫入日誌中擷取,並以新一列儲存在 Unity Catalog 管理的 Delta 表格中,每 ~ 15 秒批次處理並刷新一次。 變更歷史以開放格式儲存,任何運算引擎都能讀取。
目的資料表的結構與 Delta 變更資料摘要相同:每一列都包含 _pg_change_type、LSN、交易 ID 和時間戳記。 作業變更可成為 ETL、稽核及下游使用者的主要資料來源,而無需另外建置外部 CDC 堆疊。
應用案例
Lakebase CDF 將營運資料帶入湖屋,讓下游管線與應用程式能即時回應變化。
| 應用案例 | Description |
|---|---|
| ETL 管線 | 使用 Lakebase 作為 medallion 管線的青銅來源。 針對變更資料建立增量 Lakeflow 管線 或 Spark 結構化串流工作,並更新下游的銀表和金表。 |
| 稽核記錄 | 在 Lakebase 表格上維護完整且可查詢的每次插入、更新與刪除歷史,以符合法規與鑑識。 歷史是不可改變的三角洲。 |
| 外部系統 | 將 Lakebase 變更資料儲存為開放格式,任何引擎都能取得。 由於目的地是 Unity Catalog 中的 Delta 資料表,外部系統和非 Databricks 的讀取器可以直接存取該饋送。 |
啟用此預覽
工作區管理員必須在工作區的 預覽頁面中啟用 Lakebase 變更資料饋送預覽。
Requirements
- Lakebase: 一個運行 Postgres 16、17 或 18 的 Lakebase 專案 。
-
資料來源資料庫: 來源資料表可以存在於 Lakebase 專案中的任何單一資料庫中。 CDF 會從每個資料流擷取單一資料庫的變更;並不僅限於建立各專案時使用的
databricks_postgres資料庫。 - Unity Catalog:設定 CDF 的身分需在目標 Catalog 和結構描述上具備 USE CATALOG、USE SCHEMA 和 CREATE TABLE。 請參見 物件授權。
- 預設儲存: 預設儲存設定的目的地目錄不被支援。
- Lakebase 專案: 您在 Lakebase 專案中的 Postgres 角色需要具備 CAN MANAGE 權限。 專案擁有者預設具有可管理權限。 請參閱 管理專案權限。
- 資料類型: 參見 資料型態映射。 沒有直接 Delta 對應的型別則以 STRING 形式儲存。
Note
免費版:Databricks 免費版工作區會使用預設儲存體,供隨工作區建立的目錄使用。 若要使用 Lakebase CDF,請建立一個目錄集,其受管理的儲存位置是外部位置。
設定 Lakebase CDF
若要開始,請先在要納入資料饋送的資料表上,將 replica identity 設為 full(步驟 1),然後在 Lakebase 應用程式中啟動 CDF(步驟 2)。 你的資料會以 Delta 表格的形式出現 lb_<table_name>_history 在你選擇的 Unity 目錄和結構中。
Note
你可以從 Lakebase UI 或 API 啟動 CDF。 若要以程式管理訂閱源,請使用 Postgres REST API 和 Databricks SDK 中的 CDF 操作來建立訂閱源、檢查狀態、停用或刪除設定。 請參閱 Lakebase API 指南中的 變更資料饋源 。
步驟 1:將 Replica Identity 設為完整
Lakebase 資料表若要參與 CDF,必須設定 REPLICA IDENTITY FULL。 預設情況下,Postgres 只會在資料列更新或刪除時記錄主鍵。 設定完整識別會讓 Postgres 在預寫式日誌中同時記錄資料列變更前後的狀態,而 CDF 需要這些資訊來建立完整的變更歷程。
你可以在 Lakebase SQL Editor 或任何 Postgres 用戶端執行這些指令。
單一表格
ALTER TABLE <table_name> REPLICA IDENTITY FULL;
架構中所有現有的資料表
要在結構public 中每個現有資料表(本例中)設定副本身份,執行:
DO $$
DECLARE r record;
BEGIN
FOR r IN
SELECT table_schema, table_name
FROM information_schema.tables
WHERE table_schema = 'public'
AND table_type = 'BASE TABLE'
LOOP
EXECUTE format(
'ALTER TABLE %I.%I REPLICA IDENTITY FULL;',
r.table_schema, r.table_name
);
END LOOP;
END $$;
自動套用到未來資料表
要讓每個新建立的資料表自動收到 REPLICA IDENTITY FULL,請安裝 Postgres 事件觸發器。 它會在每次 CREATE TABLE 之後執行,並在新資料表上設定識別欄位:
CREATE OR REPLACE FUNCTION public.set_full_replica_identity()
RETURNS event_trigger
LANGUAGE plpgsql
AS $$
DECLARE
obj record;
BEGIN
FOR obj IN
SELECT * FROM pg_event_trigger_ddl_commands()
WHERE command_tag = 'CREATE TABLE'
LOOP
EXECUTE format(
'ALTER TABLE %s REPLICA IDENTITY FULL;',
obj.object_identity
);
END LOOP;
END $$;
CREATE EVENT TRIGGER set_full_replica_identity_on_create
ON ddl_command_end
WHEN TAG IN ('CREATE TABLE')
EXECUTE FUNCTION public.set_full_replica_identity();
將事件觸發器與上一個分頁中的迴圈結合,以透過一次設定涵蓋現有及未來新增的資料表。
檢查哪些資料表有複製身份設定
要查看結構中哪些資料表已設定副本身份,執行:
SELECT n.nspname AS table_schema,
c.relname AS table_name,
CASE c.relreplident
WHEN 'd' THEN 'default'
WHEN 'n' THEN 'nothing'
WHEN 'f' THEN 'full'
WHEN 'i' THEN 'index'
END AS replica_identity
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind = 'r'
AND n.nspname = 'public'
ORDER BY n.nspname, c.relname;
只有帶有 replica_identity = 'full' 的資料列才可進行 CDF。
API:REPLICA IDENTITY FULL 是標準的 Postgres DDL。 請參閱 PostgreSQL ALTER TABLE 參考文獻。
步驟 2:啟動變更資料串流
Lakebase CDF 是在結構描述層級進行設定的。 饋送開始後,來源綱要中所有目前與未來的資料表都會被包含。 若要稍後新增表格,請參見 將表格新增至執行中的摘要。
- 在你的Azure Databricks工作區,從應用程式切換器(右上角)開啟 Lakebase Postgres。
- 選擇你的 Lakebase 專案和你想使用的分支(例如 生產 或 主分支)。
- 按一下頂端導覽路徑中的分支名稱,開啟 分支總覽,然後按一下 Lakebase CDF 分頁。
- 按一下 [開始]。
- 在設定對話框中:
- 資料庫: 選擇來源 Postgres 資料庫。 你可以選擇專案中的任何資料庫,即使專案只有一個資料庫,你也必須選擇其中一個。
- 架構: 選擇來源 Postgres 架構。
- 至目錄: 選擇目標 Unity Catalog 目錄。
- 架構: 選擇目標的 Unity 目錄架構。
- 點擊 開始 播放。
表格在目的地 lb_<table_name>_history中以 的形式出現。 要找到它們,請在側邊欄開啟 目錄 ,然後進入目的地目錄和結構,然後打開 表格 分頁。
Lakebase CDF 分頁有兩個子分頁,兩者皆僅用於監控(此處無法啟用或停用單一資料表的 CDF,因為資料流始終涵蓋整個綱要範圍):
- 圖式: 列出每個來源結構、其目的目錄與 Unity 目錄中的結構,以及狀態。
-
表格: 列出每個來源資料表、其目的
lb_<table_name>_history資料表、狀態(Streaming或 或Snapshotting)、 已提交的 LSN (資料流寫入 Delta 的距離,顯示-為仍在初始快照階段)及 最後更新 (表最後一次接收變更的時間)。
你也可以透過 Lakebase SQL 編輯器執行這個來檢查 Postgres 的 feed 狀態:
SELECT * FROM wal2delta.tables;
結果包括 table_oid、status(STREAMING 或 SNAPSHOTTING)、committed_lsn 以及每個表格的 last_write_time。
Important
什麼是 wal2delta? Lakebase CDF 由 wal2delta Postgres 擴充套件驅動,該擴充功能運行於 Lakebase 運算中。 它利用邏輯解碼捕捉預先寫入日誌(WAL)變更,並將其寫入 Unity 目錄中的 Delta 表格。
API: 若要以程式方式取得資料串流設定及每表狀態,請參閱 Lakebase API 指南中的 「變更資料饋送 」。
將表格加入動態摘要
由於 CDF 的範圍限定於結構描述,你不會重新設定資料流來擷取新的資料表。 當你將資料表加入已有作用中摘要的結構時,CDF 會自動包含該資料表。 要新增一個表格:
- 在來源綱要中建立資料表,或使用現有的資料表。
- 將其 replica identity 設為 full:
ALTER TABLE <table_name> REPLICA IDENTITY FULL;。 請參見 步驟 1:完整設定複本身份。 - 確保表格至少有一列。 只有當資料表包含資料時才會同步。
表格隨即出現在目標位置 lb_<table_name>_history,且不會有其他動作。 可在 Lakebase CDF 分頁的 表格 子頁籤追蹤進度,或透過執行 SELECT * FROM wal2delta.tables;。
Warning
請勿再次點擊開始以新增資料表。 一個綱要只能有一個摘要來源。 如果你重新開啟已存在的 開始 對話方塊,該架構的資料表會顯示 衝突 狀態,警告有命名衝突的資料表會加上數字後綴。 開始對話框是用來設定尚未有 feed 的結構描述上的 feed,不是用來將資料表加入正在執行的 feed。 你加入設定結構描述的資料表會自動被擷取。
目的資料表結構
CDF 會為每個來源資料表寫出一個 Delta 表格,並命名 lb_<table_name>_history 於你的目的目錄和結構中。 除了來源欄位外,每一列還包含以下系統欄位:
| 資料行 | 類型 | Description |
|---|---|---|
_pg_change_type |
TEXT | 操作類型:insert、、deleteupdate_preimage、或update_postimage。 |
_pg_lsn |
BIGINT | Postgres 日誌序號。 |
_pg_xid |
INTEGER | Postgres 交易 ID。 |
_timestamp |
TIMESTAMP | 變更處理時間戳記(不含時區)。 |
_sort_by |
BIGINT | 單調排序鍵用於排序所有變更。 |
常見的變化模式
-
初步快照: CDF 首次在現有 Lakebase 資料表上執行時,每一列都寫成
_pg_change_type = 'insert'。 -
最新消息: 更新會產生兩列:一列為
_pg_change_type = 'update_preimage'(舊列),一列為_pg_change_type = 'update_postimage'(新列)。 -
刪除: 刪除作業會產生一個含有
_pg_change_type = 'delete'的資料列。
這些變更事件與 Delta 變更資料串流相同,因此下游模式相同。
操作行為
-
名稱衝突:如果兩個來源資料表會對應到相同的目的地名稱(例如,
sales.users和marketing.users都對應到lb_users_history),CDF 會將第一個寫入lb_users_history,並自動為第二個加上後綴,使其成為lb_users_history_1。 你可以在 Unity 目錄中重新命名任一目標資料表,訂閱串流仍能正常運作。 - 架構層級範圍: 當你在 Lakebase 架構上啟動 CDF 時,該架構中所有現有和未來的資料表都會被包含在內。 只有當資料表包含資料時才會同步,因此該資料表必須至少有一列才能出現在目的地。
- 丟棄的原始碼表: 如果你在 Lakebase 丟棄一個資料表,Unity 目錄中的目的地 Delta 資料表會被保留。
建造下游管線
Lakebase CDF 專為能因應營運變化而做出反應的下游管線設計。 下方的圖案展示了三種食用飼料的方式,依序從簡單到彈性排列。
舉個例子。 電子商務應用程式會在 Postgres orders 資料表中記錄訂單,而每一列都包含 item_id 和 quantity。 物流團隊需要即時庫存等級。 使用 CDF 時,對 orders 的每一次變更都會儲存在 Unity Catalog 中的 lb_orders_history Delta 資料表內。 下游管線會讀取該變更摘要,並在訂單建立、編輯或取消時更新 inventory_levels 資料表。
以具體化視圖計算當前庫存
最簡單的模式是對歷史資料表建立 SQL 實體化視圖 。 隨著新的變更事件到來,MV 會以增量方式更新,而下游消費者則像查詢其他資料表一樣查詢它。
CREATE MATERIALIZED VIEW inventory_levels AS
SELECT
item_id,
SUM(
CASE
-- New orders (and the "new half" of updates) decrement inventory
WHEN _pg_change_type IN ('insert', 'update_postimage') THEN -quantity
-- Cancellations (and the "old half" of updates) restore inventory
WHEN _pg_change_type IN ('delete', 'update_preimage') THEN quantity
ELSE 0
END
) AS current_inventory,
MAX(_timestamp) AS last_transaction_ts,
MAX(_pg_lsn) AS last_lsn
FROM lb_orders_history
GROUP BY item_id;
每次更新產生的兩列會互相抵銷,唯獨淨變動例外,因此在訂單編輯時,運算總和會保持正確。
Spark 宣告式管線的串流變更
若要使用結構化的勳章式架構,請使用 Lakeflow pipelines 來宣告銅、銀和金層資料表。 Lakeflow 管線會將它們作為相互串接的管線執行,並為你處理檢查點和相依性管理。
import dlt
from pyspark.sql import functions as F
@dlt.table
def inventory_adjustments():
return (
spark.readStream.table("<catalog>.<schema>.lb_orders_history")
.withColumn(
"delta",
F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
.when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
.otherwise(0),
)
.select("item_id", "delta", "_timestamp")
)
@dlt.expect_or_drop("non_negative_stock", "on_hand >= 0")
@dlt.table
def inventory_levels():
return (
spark.read.table("LIVE.inventory_adjustments")
.groupBy("item_id")
.agg(F.sum("delta").alias("on_hand"))
)
inventory_adjustments 以 lb_orders_history 增量方式讀取, readStream 並產生每個事件的 delta。
inventory_levels 依據 item_id 彙總以計算目前庫存。 預期會捨棄那些會使庫存變成負值的資料列,這表示上游有錯誤。
完整的端到端導覽,請參考 教學:使用變更資料擷取建立 ETL 管線。
使用 Spark 結構化串流進行自訂處理
當你需要完全控制權——例如自訂合併、副作用或多重匯入時——直接用 Spark Structured Streaming 讀取歷史表,然後用 foreachBatch 來寫入目的地。
from pyspark.sql import functions as F
from delta.tables import DeltaTable
def update_inventory(batch_df, batch_id):
deltas = (
batch_df
.withColumn(
"delta",
F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
.when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
.otherwise(0),
)
.groupBy("item_id")
.agg(F.sum("delta").alias("delta"))
)
target = DeltaTable.forName(spark, "<catalog>.<schema>.inventory_levels")
(target.alias("t")
.merge(deltas.alias("s"), "t.item_id = s.item_id")
.whenMatchedUpdate(set={"on_hand": F.expr("t.on_hand + s.delta")})
.whenNotMatchedInsert(values={"item_id": "s.item_id", "on_hand": "s.delta"})
.execute())
(spark.readStream.table("<catalog>.<schema>.lb_orders_history")
.writeStream
.foreachBatch(update_inventory)
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/checkpoints/inventory_levels")
.start())
每個微批次會將變更事件 item_id 彙總為 ,並將淨 delta 合併為 inventory_levels。
採漸進式設計。 每個 lb_<table_name>_history 資料表都是僅可附加的 Delta 資料表。 每次來源變更都會以新列記錄,並 _pg_change_type 標記操作內容。 Databricks 的 SQL 實體化檢視、Lakeflow pipelines 流程和 Spark Structured Streaming 工作都會從 Delta 交易日誌中逐步處理新資料列,因此下游管線的工作只會與變更的量成比例。 你不需要在歷史資料表啟用 Delta 變更資料匯送 ,因為變更語意已經編碼在資料列中。
數據類型映射
CDF 支援大多數標準 PostgreSQL 原語型別。 沒有直接 Delta 對應的型別則以 STRING 形式儲存。
| PostgreSQL 類型 | Azure Databricks Delta 類型 | 註釋 |
|---|---|---|
| BOOLEAN | BOOLEAN | |
| INT、SMALLINT、BIGINT | INT、SMALLINT、BIGINT | |
| 文字,瓦查爾,查爾 | STRING | |
| JSONB | STRING | 以 JSON 字串形式儲存。 |
| ENUM | STRING | 以列舉標籤形式儲存。 |
| 數值 / 十進位 | 十進位或字串 | 盡可能使用來源精度/比例。 對不相容的精度/縮放值進行無損縮放。 當精度超過 38 或精度/尺度未定義(無界 NUMERIC)時,會退回到 STRING。 所有 NUMERIC/DECIMAL 欄位皆可為空值,因為 NaN 值會對應為 NULL。 請參見 PostgreSQL 的數值型別。 |
| DATE | DATE | |
| TIMESTAMP | TIMESTAMP_NTZ | |
| 時間戳記 | TIMESTAMP | |
| 浮動,雙倍 | 浮動,雙倍 |
以字串形式儲存的類型:
-
地理/幾何(PostGIS): 來自 PostGIS 擴充的類型(例如,
geometry,geography)。 -
向量(pgvector):來自 pgvector 擴充功能的
vector類型。 -
複合材料/結構類型: 自訂型別定義為
CREATE TYPE ... AS (field_name type, ...)。 這些是帶有命名欄位的類列型態。 -
地圖: 類似映射的鍵值型別,例如 hstore (來自擴充套件
hstore)。 Postgres 沒有內建的地圖類型。hstore是將鍵值對儲存在欄位中的常見方式。
管理結構變更
- 在 Postgres 中重新命名資料表(例如,
ALTER TABLE users RENAME TO customers)可以讓資料流繼續。 目的地 Delta 表名稱不會改變——它保持lb_users_history。 - 結構變更(新增欄位、刪除欄位或更改欄位資料型態)會觸發受影響資料表的重新快照。 CDF 會從 Postgres 重新讀取整個資料表,並將其重寫為目的地的 Delta 資料表。
停用 Lakebase CDF
停用 CDF 會停止此專案中所有 Lakebase 結構描述的資料饋送。
- 在你的Azure Databricks工作區,從應用程式切換器(右上角)開啟 Lakebase Postgres。
- 選擇你的 Lakebase 專案以及你設定 CDF 的分支。
- 按一下頂端導覽路徑中的分支名稱,開啟 分支總覽,然後按一下 Lakebase CDF 分頁。
- 點擊 停用。 在確認對話方塊中,查看指出變更將停止流向 Delta 資料表的警告,然後再次按一下 停用 以確認。
停用 CDF 不會重新啟動你的運算。
API: 若要以程式方式停用或刪除串流設定,請參考 Lakebase API 指南中的 「變更資料饋送 」。
限制與故障排除
您可以在 Lakebase CDF 標籤中查看每表狀態(快照、跳過或串流),或在 Lakebase 執行此功能:
SELECT * FROM wal2delta.tables;
表格不會出現在動態中常見的原因:
-
REPLICA IDENTITY FULL未設定: 針對該表格執行ALTER TABLE <table_name> REPLICA IDENTITY FULL;。 請參見 步驟 1:完整設定複本身份。 - 分割資料表: 不支援 Lakebase 分割表格。 包含分割資料表的結構會導致這些資料表失敗。
- 空資料表: 只有當資料表包含資料時才會同步。 沒有列的資料表要等到至少有一列存在後,才會出現在摘要中。
Warning
請勿以以下方式修改目的 lb_<table_name>_history Delta 表格:
- 不要新增 列過濾器或欄位遮罩。 一旦你對目標資料表套用列過濾器或欄位遮罩,CDF 就會停止寫入。
- 請勿在目的地資料表啟用 Delta Lake 變更資料串流 。 它會破壞
ALTER TABLECDF 在來源結構變更時執行的重新快照功能。
Note
目的儲存上的私有端點: 當你的目標 Unity 目錄目錄的託管儲存只能透過私人端點存取時,Lakebase CDF 不被支援。 例如,AWS PrivateLink 介面端點,或是關閉公共網路存取儲存帳號的 Azure 私有端點。 作為變通方法,設定一個管理儲存空間可公開存取的目錄,並以該目錄作為你的 CDF 目的地。
下一步
- 用 Spark 宣告式管線建立增量 ETL。 完整攻略請參考 教學:使用變更資料擷取建立 ETL 管線 。
- 用 Databricks SQL查詢青銅層。 請參閱「 使用 Databricks SQL 開始資料倉儲」。
- 使用 時間旅行查詢 稽核目標 Delta 資料表上的 歷史記錄。