Lakebase 變更資料串流

Note

Lakebase 變更資料饋送功能目前處於 公開預覽階段。

什麼是 Lakebase 變更資料串流?

Lakebase 引入原生變更資料饋送(CDF),解鎖您的營運資料,用於下游的管線、模型與應用程式。 Lakebase Postgres 表格上的每次插入、更新和刪除都會從預先寫入日誌中擷取,並以新一列儲存在 Unity Catalog 管理的 Delta 表格中,每 ~ 15 秒批次處理並刷新一次。 變更歷史以開放格式儲存,任何運算引擎都能讀取。

目的資料表的結構與 Delta 變更資料摘要相同:每一列都包含 _pg_change_type、LSN、交易 ID 和時間戳記。 作業變更可成為 ETL、稽核及下游使用者的主要資料來源,而無需另外建置外部 CDC 堆疊。

Lakebase CDF 資料從 Postgres 經由 wal2delta 流向 Unity 目錄中的 Delta 表格。

應用案例

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 是在結構描述層級進行設定的。 饋送開始後,來源綱要中所有目前與未來的資料表都會被包含。 若要稍後新增表格,請參見 將表格新增至執行中的摘要。

  1. 在你的Azure Databricks工作區,從應用程式切換器(右上角)開啟 Lakebase Postgres。
  2. 選擇你的 Lakebase 專案和你想使用的分支(例如 生產 或 主分支)。
  3. 按一下頂端導覽路徑中的分支名稱,開啟 分支總覽,然後按一下 Lakebase CDF 分頁。
  4. 按一下 [開始]。
  5. 在設定對話框中:
    • 資料庫: 選擇來源 Postgres 資料庫。 你可以選擇專案中的任何資料庫,即使專案只有一個資料庫,你也必須選擇其中一個。
    • 架構: 選擇來源 Postgres 架構。
    • 至目錄: 選擇目標 Unity Catalog 目錄。
    • 架構: 選擇目標的 Unity 目錄架構。
  6. 點擊 開始 播放。

分支總覽,Lakebase CDF 標籤顯示開始與架構設定。

表格在目的地 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 會自動包含該資料表。 要新增一個表格:

  1. 在來源綱要中建立資料表,或使用現有的資料表。
  2. 將其 replica identity 設為 full:ALTER TABLE <table_name> REPLICA IDENTITY FULL;。 請參見 步驟 1:完整設定複本身份。
  3. 確保表格至少有一列。 只有當資料表包含資料時才會同步。

表格隨即出現在目標位置 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 結構描述的資料饋送。

  1. 在你的Azure Databricks工作區,從應用程式切換器(右上角)開啟 Lakebase Postgres。
  2. 選擇你的 Lakebase 專案以及你設定 CDF 的分支。
  3. 按一下頂端導覽路徑中的分支名稱,開啟 分支總覽,然後按一下 Lakebase CDF 分頁。
  4. 點擊 停用。 在確認對話方塊中,查看指出變更將停止流向 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 TABLE CDF 在來源結構變更時執行的重新快照功能。

Note

目的儲存上的私有端點: 當你的目標 Unity 目錄目錄的託管儲存只能透過私人端點存取時,Lakebase CDF 不被支援。 例如,AWS PrivateLink 介面端點,或是關閉公共網路存取儲存帳號的 Azure 私有端點。 作為變通方法,設定一個管理儲存空間可公開存取的目錄,並以該目錄作為你的 CDF 目的地。

下一步