Lakebase-ändringsdataflöde

Note

Funktionen Ändra dataflöde i Lakebase finns i offentlig förhandsversion.

Vad är Lakebase Change Data Feed?

Lakebase introducerar en inbyggd CDF (Change Data Feed) som låser upp dina driftdata för nedströmspipelines, modeller och program. Varje infogning, uppdatering och borttagning i en Lakebase Postgres-tabell samlas in från loggen för framåtskrivning och lagras som en ny rad i en hanterad Delta-tabell i Unity Catalog, batchad och tömd var 15:e sekund. Ändringshistoriken lagras i ett öppet format som alla beräkningsmotorer kan läsa.

Måltabellerna följer samma form som Delta Change Data Feed: varje rad har ett _pg_change_type, ett LSN, ett transaktions-ID och en tidsstämpel. Ändringar i driften blir en primär källa för ETL, revision och nedströmskonsumenter – utan att behöva sätta upp en extern CDC-stack.

Lakebase CDF-dataflöde från Postgres via wal2delta till Delta-tabeller i Unity Catalog.

Användningsfall

Lakebase CDF för in operativa data i lakehouse-plattformen så att efterföljande pipelines och applikationer kan reagera på ändringar i takt med att de sker.

Användningsfall Description
ETL-pipelines Använd Lakebase som bronskälla för medaljongpipelines. Skapa inkrementella Lakeflow-pipelines eller Spark Structured Streaming-jobb baserat på ändringsflödet och uppdatera nedströmsliggande silver- och guldtabeller.
Granskningsloggar Behåll en fullständig, sökbar historik över varje infogning, uppdatering och borttagning i en Lakebase-tabell för efterlevnad och forensisk analys. Historiken är oföränderbar Delta.
Externa system Lagra Lakebase ändringsdata i ett öppet format som alla bearbetningsmotorer kan läsa. Eftersom destinationen är en Delta-tabell i Unity Catalog kan externa system och läsare utanför Databricks komma åt flödet direkt.

Aktivera den här förhandsversionen

En administratör för arbetsytan måste aktivera förhandsversionen av Lakebase Change Data Feed på arbetsytans Förhandsversioner-sida.

Requirements

  • Lakebase:Ett Lakebase-projekt som kör Postgres 16, 17 eller 18.
  • Källdatabas: Källtabeller kan finnas i vilken enskild databas som helst i ditt Lakebase-projekt. CDF fångar upp ändringar från en databas per flöde; Det är inte begränsat till databasen databricks_postgres som varje projekt skapas med.
  • Unity-katalog: Den identitet som konfigurerar CDF behöver USE CATALOG, USE SCHEMAoch CREATE TABLE i målkatalogen och schemat. Se Bevilja behörigheter för ett objekt.
  • Standardlagring: Målkataloger som konfigurerats med standardlagring stöds inte.
  • Lakebase-projekt: Din Postgres-roll kräver CAN MANAGE-behörigheter för Lakebase-projektet. Projektägare har KAN HANTERA som standard. Se Hantera projektbehörigheter.
  • Datatyper: Se Datatypsmappning. Typer utan en direkt Delta-motsvarighet lagras som STRING.

Note

Free Edition:Databricks Free Edition-arbetsytor använder standardlagring för katalogen som skapas med arbetsytan. För att använda Lakebase CDF, skapa en katalog vars hanterade lagringsplats är en extern plats.

Konfigurera Lakebase CDF

Börja med att ställa in REPLICA IDENTITY FULL på de tabeller som du vill ha i flödet (steg 1) och starta sedan CDF i Lakebase-appen (steg 2). Dina data visas som lb_<table_name>_history Delta-tabeller i Den Unity Catalog-katalog och det schema som du väljer.

Note

Du kan starta CDF från Lakebase UI eller med API:et. För att hantera ett flöde programmatiskt, använd CDF-operationerna i Postgres REST API och Databricks SDK:er för att skapa ett flöde, kontrollera dess status, inaktivera det eller ta bort en konfiguration. Se Change Data Feed i Lakebase API-guiden.

Steg 1: Ange replikaidentitet till full

För att en Lakebase-tabell ska kunna delta i CDF måste den ha REPLICA IDENTITY FULL angetts. Postgres loggar som standard endast primärnyckeln när en rad uppdateras eller tas bort. Att ställa in full identitet innebär att Postgres registrerar både radens tillstånd före och efter i write-ahead-loggen, vilket CDF behöver för att bygga upp en fullständig ändringshistorik.

Du kan köra dessa kommandon i Lakebase SQL-redigeraren eller någon Postgres-klient.

Enkel tabell

ALTER TABLE <table_name> REPLICA IDENTITY FULL;

Alla befintliga tabeller i ett schema

Om du vill ange replikidentitet för varje befintlig tabell i ett schema (public i det här exemplet) kör du:

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 $$;

Tillämpa automatiskt på framtida tabeller

Om du vill att varje nyskapad tabell automatiskt ska ta emot REPLICA IDENTITY FULLinstallerar du en Postgres-händelseutlösare. Den körs efter varje CREATE TABLE och anger identiteten i den nya tabellen:

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();

Kombinera händelseutlösaren med loopen på föregående flik för att täcka både befintliga och framtida tabeller i en konfiguration.

Kontrollera vilka tabeller som har replikidentitet angiven

Om du vill se vilka tabeller i ett schema som har replikidentiteten konfigurerad kör du:

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;

Endast rader med replica_identity = 'full' är klara för CDF.

API:REPLICA IDENTITY FULL är standard Postgres DDL. Se PostgreSQL-ALTER TABLEreferensen.

Steg 2: Starta ändringsdataflödet

Lakebase CDF konfigureras på schemanivå. När den har startats inkluderas varje aktuell och framtida tabell i källschemat i feeden.

  1. I din Azure Databricks arbetsyta öppnar du Lakebase Postgres från appväxlaren (längst upp till höger).
  2. Välj ditt Lakebase-projekt och den gren som du vill använda (till exempel produktion eller huvud).
  3. Öppna Grenöversikt genom att klicka på grennamnet i sökvägen högst upp och klicka sedan på fliken Lakebase CDF.
  4. Klicka på Start.
  5. I konfigurationsdialogrutan:
    • Databas: Välj källan Postgres-databasen. Du kan välja vilken databas som helst i projektet, och du måste välja en även om projektet bara har en databas.
    • Schemat: Välj postgres-källschemat.
    • Till katalog: Välj målkatalogen i Unity Catalog.
    • Schemat: Välj målschemat för Unity-katalogen.
  6. Klicka på Start för att starta feeden.

Grenöversikt med fliken Lakebase CDF som visar start- och schemakonfiguration.

Tabeller visas på målplatsen som lb_<table_name>_history. Om du vill hitta dem öppnar du Katalogen i sidofältet, navigerar till målkatalogen och schemat och öppnar fliken Tabeller .

Fliken Lakebase CDF har två underflikar:

Underflikar visar mappningen och framstegen per tabell.

  • Scheman: Visar en lista över varje källschema, dess målkatalog och schema i Unity Catalog och en status.
  • Tabeller: Listar varje källtabell, dess motsvarande måltabell lb_<table_name>_history, status (Streaming eller Snapshotting), Bekräftad LSN (hur långt flödet har skrivit till Delta, visas som - så länge den första ögonblicksbilden fortfarande pågår) och Senaste uppdatering (senaste gången tabellen fick ändringar).

Du kan också kontrollera flödestillståndet från Postgres genom att köra detta i Lakebase SQL-redigeraren:

SELECT * FROM wal2delta.tables;

Resultatet inkluderar table_oid, status (STREAMING eller SNAPSHOTTING), committed_lsnoch last_write_time per tabell.

Important

Vad är wal2delta? Lakebase CDF drivs av wal2delta Postgres-tillägget, som körs inuti Lakebase-beräkningen. Den använder logisk avkodning för att samla in wal-ändringar (write-ahead log) och skriver dem till Delta-tabeller i Unity Catalog.

API: För att programmatiskt hämta flödeskonfigurationen och statusen per tabell, se Ändra dataflöde i Lakebase API-guiden.

Schema för måltabell

CDF skriver en Delta-tabell per källtabell med namnet lb_<table_name>_history i målkatalogen och schemat. Förutom dina källkolumner har varje rad dessa systemkolumner:

Column Type Description
_pg_change_type Textmeddelande Åtgärdstyp: insert, delete, update_preimageeller update_postimage.
_pg_lsn BIGINT Postgres-loggsekvensnummer.
_pg_xid INTEGER Postgres transaktions-ID.
_timestamp TIMESTAMP Tidsstämpel när ändringen bearbetades (utan tidszon).
_sort_by BIGINT Monoton sorteringsnyckel som används för att ordna alla ändringar.

Vanliga ändringsmönster

  • Första ögonblicksbilden: Första gången CDF körs på en befintlig Lakebase-tabell skrivs varje befintlig rad med _pg_change_type = 'insert'.
  • Uppdateringar: En uppdatering genererar två rader: en med _pg_change_type = 'update_preimage' (gammal rad) och en med _pg_change_type = 'update_postimage' (ny rad).
  • Tar bort: En borttagning genererar en rad med _pg_change_type = 'delete'.

Det här är samma ändringshändelser som Delta Change Data Feed, så samma underordnade mönster gäller.

Driftsbeteende

  • Namngivningskollisioner: Om två källtabeller mappas till samma målnamn (till exempel sales.users och marketing.users båda mappar till lb_users_history), skriver CDF det första till lb_users_history och det automatiska suffixet det andra till lb_users_history_1. Du kan byta namn på någon av måltabellerna i Unity Catalog och feeden fortsätter att fungera.
  • Omfång på schemanivå: När du startar CDF på ett Lakebase-schema inkluderas varje aktuell och framtida tabell i schemat. Tomma tabeller hoppas över – en tabell måste ha minst en rad för att visas i destinationen.
  • Borttagna källtabeller: Om du tar bort en tabell i Lakebase bevaras mål-Delta-tabellen i Unity Catalog.

Skapa efterföljande pipelines

Lakebase CDF är utformat för nedströmspipelines som reagerar på driftsändringar. Mönstren nedan visar tre sätt att ta del av flödet, ordnade från det enklaste till det mest flexibla.

Exempelscenario. En e-handelsapp registrerar beställningar i en Postgres orders tabell, där varje rad innehåller en item_id och quantity. Logistikteamet behöver lagernivåer i realtid. Med CDF lagras varje ändring av orders i Delta-tabellen lb_orders_history i Unity Catalog. Underordnade pipelines läser ändringsflödet och uppdaterar en inventory_levels tabell när en beställning görs, redigeras eller avbryts.

Beräkna aktuell inventering med en materialiserad vy

Det enklaste mönstret är en SQL-materialiserad vy över historiktabellen. MV uppdateras stegvis när nya ändringshändelser anländer, och nedströmskonsumenter gör frågor mot den precis som mot vilken annan tabell som helst.

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;

De två rader som skapas för varje uppdatering avbryter varandra förutom nettoändringen, så den löpande summan förblir korrekt när beställningar redigeras.

Dataströmma ändringar med Spark Declarative Pipelines

För en strukturerad medaljongarkitektur kan du använda Lakeflow-pipelines för att definiera brons-, silver- och guldtabeller. Lakeflow-pipelines kör dem som en ansluten pipeline med kontrollpunkter och beroendehantering som hanteras åt dig.

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 läser lb_orders_history stegvis med readStream och genererar ett delta per händelse. inventory_levels aggregerar efter item_id för att beräkna aktuellt lager. Förväntningen utelämnar rader som skulle göra lagersaldot negativt, vilket signalerar ett fel i ett tidigare led.

En fullständig genomgång från slutpunkt till slutpunkt finns i Självstudie: Skapa en ETL-pipeline med hjälp av insamling av ändringsdata.

Anpassad bearbetning med Spark Structured Streaming

När du behöver fullständig kontroll – till exempel anpassade sammanslagningar, biverkningar eller flera mottagare – läser du historiktabellen direkt med Spark Structured Streaming och använder foreachBatch för att skriva till ditt mål.

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())

Varje mikrobatch aggregerar ändringshändelserna efter item_id och sammanfogar nettodeltorna till inventory_levels.

Utformad för att vara inkrementell. Varje lb_<table_name>_history tabell är en Delta-tabell där endast tillägg tillåts. Varje källändring registreras som en ny rad där _pg_change_type åtgärden markeras. Materialiserade vyer i Databricks SQL, flöden i Lakeflow pipelines och Spark Structured Streaming-jobb bearbetar alla nya rader inkrementellt från Deltas transaktionslogg, så efterföljande pipelines bara behöver utföra arbete i proportion till det som har ändrats. Du behöver inte aktivera Delta Change Data Feed i historiktabellen eftersom ändringssemantik redan är kodade i raddata.

Datatypkartläggning

CDF stöder de flesta vanliga PostgreSQL-primitiva typer. Typer utan en direkt Delta-motsvarighet lagras som STRING.

PostgreSQL-typ Azure Databricks Delta-typ Notes
Boolean Boolean
INT, SMALLINT, BIGINT INT, SMALLINT, BIGINT
TEXT, VARCHAR, CHAR STRING
JSONB STRING Lagras som en JSON-sträng.
ENUM STRING Lagras som enum-etikett.
NUMERISK/DECIMAL DECIMAL ELLER STRÄNG Använder källans precision/skala om möjligt. Utför förlustfri omskalning för inkompatibla precisions-/skalningsvärden. Återgår till STRING när precisionen överskrider 38 eller när precision/skala är odefinierad (obundna NUMERISKA). Alla NUMERISKA/DECIMAL-kolumner är null eftersom NaN-värden mappas till NULL. Se Numeriska postgreSQL-typer.
DATE DATE
TIMESTAMP TIMESTAMP_NTZ
TIDSSTÄMPEL TIMESTAMP
FLYT, DUBBEL FLYT, DUBBEL

Typer som lagras som STRING:

  • Geografi/geometri (PostGIS): Typer från PostGIS-tillägget (till exempel geometry, geography).
  • Vektor (pgvector): Typen vector från pgvector-tillägget.
  • Sammansatta/struct-typer: Anpassade typer som definierats med CREATE TYPE ... AS (field_name type, ...). Det här är radliknande typer med namngivna fält.
  • Map: Nyckel-värde-typer av map-typ, till exempel hstore (från tillägget hstore). Postgres har ingen inbyggd karttyp. hstore är det vanliga sättet att lagra nyckel/värde-par i en kolumn.

Hantera schemaändringar

  • Genom att byta namn på en tabell i Postgres (till exempel ALTER TABLE users RENAME TO customers) kan flödet fortsätta. Målnamnet för Delta-tabellen ändras inte – det förblir lb_users_history.
  • Schemaändringar (lägga till en kolumn, släppa en kolumn eller ändra en kolumns datatyp) utlöser en ny ögonblicksbild av den berörda tabellen. CDF läser om hela tabellen från Postgres och skriver om den till deltatabellen.

Inaktivera Lakebase CDF

Om du inaktiverar CDF stoppas flödet för alla Lakebase-scheman i projektet.

  1. I din Azure Databricks arbetsyta öppnar du Lakebase Postgres från appväxlaren (längst upp till höger).
  2. Välj ditt Lakebase-projekt och den gren där du konfigurerade CDF.
  3. Öppna Grenöversikt genom att klicka på grennamnet i sökvägen högst upp och klicka sedan på fliken Lakebase CDF.
  4. Klicka på Inaktivera. I bekräftelsedialogrutan granskar du varningen om att ändringarna slutar flöda till Delta-tabeller och klickar sedan på Inaktivera igen för att bekräfta.

Om du inaktiverar CDF startas inte din beräkning om.

API: För att inaktivera eller ta bort en feedkonfiguration programmatiskt, se Change Data Feed i Lakebase API-guiden.

Begränsningar och felsökning

Du kan se status för varje tabell (ögonblicksbild, överhoppad eller strömmande) på fliken Lakebase CDF, eller genom att köra följande i Lakebase:

SELECT * FROM wal2delta.tables;

Vanliga orsaker till att en tabell inte visas i feeden:

  • REPLICA IDENTITY FULL inte inställt: Kör ALTER TABLE <table_name> REPLICA IDENTITY FULL; för tabellen. Se Steg 1: Ställ in replikidentitet till full.
  • Partitionerade tabeller: Lakebase-partitionerade tabeller stöds inte. Ett schema som innehåller partitionerade tabeller gör att tabellerna misslyckas.
  • Tomma tabeller: En tabell utan rader hoppas över tills den innehåller minst en rad.

Varning

Ändra inte destinations-Delta-tabellerna lb_<table_name>_history på följande sätt:

  • Lägg inte till radfilter eller kolumnmasker. CDF slutar skriva till en destinationstabell när du applicerar ett radfilter eller kolumnmask på den.
  • Aktivera inte Delta Lake-ändringsdataflödet i en destinationstabell. Det gör att den ALTER TABLE nya ögonblicksbild som CDF kör när källschemat ändras inte fungerar.

Note

Privat slutpunkt på destinationslagring: Lakebase CDF stöds inte när det hanterade lagret för din destinationskatalog i Unity Catalog endast är tillgängligt via en privat endpoint. Exempel inkluderar en AWS PrivateLink-gränssnittsendpoint, eller en Azure privat endpoint med offentlig nätverksåtkomst till lagringskontot inaktiverat. Som en lösning konfigurerar du en katalog vars hanterade lagring kan nås offentligt och använder katalogen som CDF-mål.

Nästa steg