Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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.
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_postgressom 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.
- I din Azure Databricks arbetsyta öppnar du Lakebase Postgres från appväxlaren (längst upp till höger).
- Välj ditt Lakebase-projekt och den gren som du vill använda (till exempel produktion eller huvud).
- Öppna Grenöversikt genom att klicka på grennamnet i sökvägen högst upp och klicka sedan på fliken Lakebase CDF.
- Klicka på Start.
- 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.
- Klicka på Start för att starta feeden.
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:
- 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 (StreamingellerSnapshotting), 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.usersochmarketing.usersbåda mappar tilllb_users_history), skriver CDF det första tilllb_users_historyoch det automatiska suffixet det andra tilllb_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
vectorfrå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örblirlb_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.
- I din Azure Databricks arbetsyta öppnar du Lakebase Postgres från appväxlaren (längst upp till höger).
- Välj ditt Lakebase-projekt och den gren där du konfigurerade CDF.
- Öppna Grenöversikt genom att klicka på grennamnet i sökvägen högst upp och klicka sedan på fliken Lakebase CDF.
- 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 FULLinte inställt: KörALTER 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 TABLEnya ö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
- Bygg inkrementell ETL med Spark Declarative Pipelines. Se Självstudie: Skapa en ETL-pipeline med ändringsdatainsamling för en fullständig genomgång.
- Fråga bronsskiktet med Databricks SQL. Se Komma igång med datalagerhantering med Databricks SQL.
- Granska historiken med time travel-frågor på mål-Delta-tabellerna.