Avancerade AUTO CDC-ämnen

Utöver de grundläggande AUTO CDC API:erna och AUTO CDC FROM SNAPSHOT API:erna kan du köra DML på måltabeller, läsa ändringsdataflöden från CDC-mål, övervaka bearbetningsmått, tillämpa partiella uppdateringar och spåra ändringar med bitemporal storage. En introduktion till API:erna finns i AUTO CDCAPI:er för AUTOMATISK CDC: Förenkla datainsamling av ändringar med pipelines.

Lägga till, ändra eller ta bort data i en direktuppspelningstabell

Om din pipeline publicerar tabeller till Unity Catalog kan du använda DML-instruktioner ( datamanipuleringsspråk ), inklusive infognings-, uppdaterings-, borttagnings- och sammanslagningsinstruktioner, för att ändra målströmningstabeller som skapats av AUTO CDC ... INTO -instruktioner.

Anmärkning

  • DML-instruktioner som ändrar tabellschemat för en strömmande tabell stöds inte. Se till att DML-uttrycken inte försöker utveckla tabellschemat.
  • DML-instruktioner som uppdaterar en strömmande tabell kan endast köras i ett delat Unity Catalog-kluster eller ett SQL-lager med Databricks Runtime 13.3 LTS och senare.
  • Eftersom direktuppspelning kräver tilläggsdatakällor anger du flaggan skipChangeCommits när du läser källströmningstabellen om bearbetningen kräver strömning från en källströmningstabell med ändringar (till exempel av DML-instruktioner). När skipChangeCommits anges ignoreras transaktioner som tar bort eller ändrar poster i källtabellen. Om din bearbetning inte kräver en strömmande tabell kan du använda en materialiserad vy (som inte har append-only-begränsningen) som måltabell.

Eftersom pipelinen använder en angiven SEQUENCE BY kolumn och sprider lämpliga sekvenseringsvärden till måltabellens __START_AT kolumner och __END_AT (för SCD-typ 2) måste du se till att DML-uttryck använder giltiga värden för dessa kolumner för att upprätthålla rätt ordning på posterna. Se Hur AUTO CDC fungerar.

Mer information om hur du använder DML-instruktioner med strömmande tabeller finns i Lägga till, ändra eller ta bort data i en strömmande tabell.

I följande exempel infogas en aktiv post med en startsekvens på 5:

INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);

Tips/Råd

Om du behöver byta namn på kolumnerna och __START_AT i måltabellen __END_AT för SCD Typ 2 (till exempel för att matcha schemakraven för nedströms) skapar du en vy över måltabellen:

CREATE VIEW my_employees_view AS
SELECT
  *,
  __START_AT AS valid_from,
  __END_AT AS valid_to
FROM my_scd2_target_table;

Läsa ett ändringsdataflöde från en AUTO CDC-måltabell

I Databricks Runtime 15.2 och senare kan du läsa ett ändringsdataflöde från en strömmande tabell som utgör måltavlan för AUTO CDC eller AUTO CDC FROM SNAPSHOT querys på samma sätt som du läser ett ändringsdataflöde från andra Delta-tabeller. Följande krävs för att läsa ändringsdataflödet från en målströmningstabell:

  • Målsströmningstabellen måste publiceras till Unity Catalog. Se Använd Unity Catalog med pipelines.
  • Om du vill läsa ändringsdataflödet från målströmningstabellen måste du använda Databricks Runtime 15.2 eller senare. Om du vill läsa ändringsdataflödet i en annan pipeline måste pipelinen konfigureras för att använda Databricks Runtime 15.2 eller senare.

Du läser ändringsdataflödet från en målströmningstabell som skapades i en Lakeflow-pipeline på samma sätt som när du läste ett ändringsdataflöde från andra Delta-tabeller. Mer information om hur du använder funktionen deltaändringsdataflöde, inklusive exempel i Python och SQL, finns i Använda ändringsdataflöde på Azure Databricks.

Anmärkning

Posten för ändringsdataflöde innehåller metadata som identifierar typen av ändringshändelse. När en post uppdateras i en tabell, innehåller medatadatan för associerade ändringsposter vanligtvis värden som är inställda på _change_type och update_preimage händelser.

Värdena skiljer sig dock _change_type om uppdateringar görs i måluppspelningstabellen som inkluderar ändring av primärnyckelvärden. När ändringar inkluderar uppdateringar av primära nycklar anges metadatafälten _change_type till insert och delete händelser. Ändringar i primära nycklar kan ske när manuella uppdateringar görs i ett av nyckelfälten med en UPDATE eller MERGE instruktion eller, för SCD-tabeller av typ 2, när fältet __start_at ändras för att återspegla ett tidigare startsekvensvärde.

Frågan AUTO CDC avgör de primära nyckelvärdena, som skiljer sig åt för SCD-typ 1- och SCD-typ 2-bearbetning:

SCD-typ Primärnyckel
SCD-typ 1 och Python-gränssnittet för pipelines Primärnyckeln är värdet för parametern keys i create_auto_cdc_flow() funktionen. För SQL-gränssnittet är den primära nyckeln de kolumner som definieras av KEYS -satsen i -instruktionen AUTO CDC ... INTO .
SCD-typ 2 Den primära nyckeln är parametern keys eller KEYS satsen plus returvärdet från coalesce(__START_AT, __END_AT) åtgärden, där __START_AT och __END_AT är motsvarande kolumner från måluppspelningstabellen. Detta använder __START_AT när det är tillgängligt och __END_AT när __START_AT är null (till exempel den första posten).

Läs ett ändringsdataflöde från en materialiserad vy

Important

Den här funktionen finns i Beta.

Du kan läsa ett ändringsdataflöde från en materialiserad vy skapad i en Lakeflow-pipeline eller i Databricks SQL. Använd detta för att replikera materialiserade vyändringar till destinationer utanför Azure Databricks, eller för att föra en historik över materialiserade vyändringar för revision och rapportering.

Materialiserade vyer använder automatiskt ändringsdataflöde, så du aktiverar inte själva ändringsdataflödet. Istället aktiverar du ändringsdata-flödet på varje materialiserad vy där du behöver det genom att uppfylla följande krav. Se Automatisk ändringsdataflöde.

  • För att läsa flödet för ändringsdata måste du använda Databricks Runtime 18 LTS eller senare, med klassisk beräkning, serverlös beräkning eller Databricks SQL.

  • Den materialiserade vyn, pipelinen som skapar den eller pipelinen som läser den måste använda PREVIEWkanalen.

  • Den materialiserade vyn måste ha radspårning aktiverad. Materialiserade vyer på serverlös beräkning har radspårning aktiverad som standard. Se Radspårning i Azure Databricks. För att kontrollera om radspårning är aktiverat i en materialiserad vy, kör:

    SHOW TBLPROPERTIES my_mv ('delta.enableRowTracking');
    
  • För att läsa ändringsdataflödet från en materialiserad vy, aktivera den externa metadataflaggan på pipelinen eller den materialiserade vyn. För instruktioner, se Hur man aktiverar åtkomst för en datamängd.

Du läser ändringsdataflödet från en materialiserad vy på samma sätt som från andra Delta-tabeller, med funktionen table_changes() , en strömmande läsning eller alternativet readChangeFeed . För syntax och exempel i SQL och Python, se Use change data feed on Azure Databricks.

Du kan läsa ett ändringsdataflöde för en materialiserad vy i en materialiserad vy eller strömmande tabell i Databricks SQL:

CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT * FROM STREAM my_mv WITH (readChangeFeed=true)

Limitations

Utöver begränsningarna för automatiska ändringsdataflöden gäller följande när du läser ett ändringsdataflöde från en materialiserad vy:

  • Ändringsdataflödet inkluderar oförändrade rader när den materialiserade vyn är helt omskriven, och den konsoliderar inte flera uppdateringar till samma rad till en enda händelse. För att filtrera bort dessa, aggregera ändringsdataflödet genom att gruppera på alla kolumner för att hitta insättningar och borttagningar som delar samma radvärden.
  • Endast Azure Databricks kan söka ändringsdataflödet efter en materialiserad vy. Externa Delta Lake- och Iceberg-klienter kan inte det.
  • Inom Lakeflow-pipelines kan du läsa ändringsdataflödet för en materialiserad vy endast från en annan pipeline, och den pipelinen måste använda kanalen PREVIEW. Att läsa ändringsdataflödet för en materialiserad vy i samma pipeline som skapar den stöds inte.
  • Du kan inte skapa ett vektorsökningsindex från en materialiserad vy.

Hämta data om poster som bearbetas av en CDC-fråga i pipelines

Anmärkning

Följande mått registreras endast av AUTO CDC frågor och inte av AUTO CDC FROM SNAPSHOT frågor.

Följande mått samlas in av AUTO CDC frågor:

  • num_upserted_rows: Antalet utdatarader som har infogats till datamängden under en uppdatering.
  • num_deleted_rows: Antalet befintliga utdatarader som tagits bort från datauppsättningen under en uppdatering.

Måttet num_output_rows, som är utdata för icke-CDC-flöden, samlas inte in för AUTO CDC sökfrågor.

Tillämpa partiella uppdateringar

När en källa bara skickar de kolumner som har ändrats AUTO CDC måste skilja mellan en kolumn som saknas från en ändringspost, vilket bör lämna målvärdet oförändrat och en kolumn som uttryckligen är inställd på null, som ska skriva över målvärdet med null. Som standard IGNORE NULL UPDATES behandlar alla null som en "uppdatera inte"-markör, så det kan inte tillämpa en explicit null. Lös den här tvetydigheten genom att välja någon av följande tre metoder:

Method När det bör användas Behavior
IGNORE NULL UPDATES ON columnList En liten, fast uppsättning kolumner bör ignorera null värden, medan alla andra kolumner använder explicita null värden. De listade kolumnerna behåller sitt befintliga målvärde när det inkommande värdet är null. Alla andra kolumner använder explicita null värden.
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList) De flesta kolumner bör ignorera null värden, och endast ett fåtal bör använda explicita null värden. Kolumnerna i listan använder explicita null värden. Alla andra kolumner behåller sitt befintliga målvärde när det inkommande värdet är null.
COLUMNS TO UPDATE Varje ändringspost uppdaterar en annan uppsättning kolumner, eller så ändras uppsättningen med uppdateringsbara kolumner över tid. En källkolumn namnger de kolumner som ska uppdateras för varje ändringspost. De listade kolumnerna skrivs från källan, inklusive explicita null värden. Kolumner som inte visas behåller sitt befintliga målvärde.

COLUMNS TO UPDATE kan inte kombineras med IGNORE NULL UPDATESoch stöds inte för bitemporal-tabeller.

Välj som tumregel COLUMNS TO UPDATE när producenten vet vilka kolumner som har ändrats i varje post och kan ange den informationen i en källkolumn, till exempel när flera producenter skriver till samma datakälla eller när uppsättningen uppdateringsbara kolumner växer över tid. Välj IGNORE NULL UPDATES ON när pipelineägaren känner till den fasta uppsättningen med uppdateringsbara kolumner i förväg och föredrar att styra dem i pipelinekoden.

I följande exempel används en källkolumn med namnet columnsToUpdate för att styra vilka kolumner varje ändringspost uppdaterar, inklusive kolumner som uttryckligen anges till null:

Python

from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
  target = "target",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "sequenceNum",
  stored_as_scd_type = 1,
  columns_to_update = "columnsToUpdate"
)

SQL

CREATE OR REFRESH STREAMING TABLE target;

CREATE FLOW apply_cdc AS AUTO CDC INTO
  target
FROM
  stream(cdc_source)
KEYS
  (id)
SEQUENCE BY
  sequenceNum
STORED AS
  SCD TYPE 1
COLUMNS TO UPDATE
  columnsToUpdate;

En fullständig referens för parametrarna finns i AUTO CDC INTO (pipelines) och create_auto_cdc_flow.

Bitemporal AUTO CDC

Important

Bitemporal AUTO CDC är i betaversion.

SCD Typ 1 och Typ 2 är unitemporal: de spårar ändringar i en enda tidsdimension. Bitemporal utökar SCD Typ 2-historiken för att spåra ändringar över två tidsdimensioner och skilja mellan två perspektiv:

  • Affärstid: när händelsen faktiskt inträffade.
  • Systemtid: när systemet registrerade eller matade in händelsen.

Liksom SCD Typ 2 bevarar bitemporal en fullständig historik över poster. Den lägger till en andra tidslinje så att du kan rekonstruera både vad data visade och vad systemet trodde när som helst tidigare.

En hedgefond matar till exempel in aktiedata från ett källsystem. Acme Corps aktiekurs ändras den 1 januari, men fonden matar inte in den uppdateringen förrän den 5 januari. Bitemporal AUTO CDC låter fonden svara på två distinkta frågor: vad Acme Corps faktiska aktiekurs var den 1 januari (affärstid) och vilket pris systemet trodde när fonden fattade handelsbeslut den 3 januari (systemtid). Möjligheten att skilja mellan dessa tidslinjer är användbar för granskning, regelrapportering och ekonomiskt beslutsfattande.

Om du vill aktivera bitemporal bearbetning anger du STORED AS BITEMPORAL (SQL) eller stored_as_scd_type="bitemporal" (Python), använder SEQUENCE BY för kolumnen affärstid och använder SYSTEM SEQUENCE BY för systemtidskolumnen. Måltabellen lägger till __SYSTEM_START_AT och __SYSTEM_END_AT kolumner tillsammans med SCD Typ 2 __START_AT och __END_AT kolumner. Syntaxinformation finns i AUTO CDC INTO (pipelines) eller create_auto_cdc_flow.

Bitemporal AUTO CDC-exempel

Följande exempel skapar en bitemporal måltabell från ett mindre antal syntetiska CDC-händelser. Kolumnen bt har affärstiden och st kolumnen har systemtiden.

Python

from pyspark import pipelines as dp

# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")

@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
  return spark.createDataFrame(
    [
      (1, "x10", "y10", 10, 100),
      (1, "x20", "y20", 20, 200)
    ],
    schema="id INT, x STRING, y STRING, bt INT, st INT",
  )

# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")

dp.create_auto_cdc_flow(
  target = "target_bitemporal",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "bt",
  system_sequence_by = "st",
  stored_as_scd_type = "bitemporal"
)

SQL

-- Source: synthetic CDC events
CREATE OR REFRESH STREAMING TABLE cdc_source_sql;

CREATE FLOW cdc_source_sql AS INSERT INTO ONCE
  cdc_source_sql BY NAME
SELECT * FROM VALUES
  (1, 'x10', 'y10', 10, 100),
  (1, 'x20', 'y20', 20, 200)
  AS t(id, x, y, bt, st);

-- Target: bitemporal table
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;

CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
  target_bitemporal_sql
FROM
  stream(cdc_source_sql)
KEYS
  (id)
SEQUENCE BY
  bt
SYSTEM SEQUENCE BY
  st
STORED AS
  BITEMPORAL;

Följande sekvens med ändringar visar hur en bittabell registrerar en infogning, en uppdatering, en oordnad uppdatering och en borttagning för ett enskilt företag. Sekvenseringskolumnen genererar kolumnerna __START_AT och __END_AT (affärstid) och systemsekvenseringskolumnen genererar kolumnerna __SYSTEM_START_AT och __SYSTEM_END_AT (systemtid):

Column Description
__START_AT Den verksamhetstid vid vilken denna rad blev giltig.
__END_AT Den affärstid då den här radens giltighet upphör. null om den är giltig på obestämd tid.
__SYSTEM_START_AT Systemtiden då den här radens data- och affärstidsintervall är kända för att vara sanna.
__SYSTEM_END_AT Systemtiden då den här radens data- och affärstidsintervall är kända för att vara ogiltiga. null om det är känt att det är sant på obestämd tid.

Systemet hanterar händelser som anländer i valfri ordning över båda tidslinjerna. När en händelse kommer med en tidigare affärstid eller systemtid än händelser som redan har bearbetats korrigerar systemet den berörda historiken i stället för att bara lägga till till slutet.

Ändring 1: Infoga

Företag A läggs till 2025-07-18 10:01:00 (affärstid) men matas inte in förrän 10:05:00 (systemtid).

Input:

CompanyId Datapunkt Ordningsföljd Systemsekvensering Operation
A XFv1 7/18/2025 10:01:00 7/18/2025 10:05:00 INSERT

Resultat:

CompanyId Datapunkt __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULL 7/18/2025 10:05:00 NULL

XFv1 är giltigt från och med 10:01:00 utan kända slut. Systemet fick reda på det här faktumet vid systemtiden 10:05:00, utan något känt slut.

Ändring 2: Uppdatera

Företag A uppdaterades 2025-07-18 kl. 12:15:43 (affärstid), och systemet behandlar händelsen kl. 12:20:00 (systemtid). Systemet bevarar både vad man trodde innan uppdateringen var känd och den korrigerade affärshistoriken efter att uppdateringen har matats in.

Input:

CompanyId Datapunkt Ordningsföljd Systemsekvensering Operation
A XFv2 7/18/2025 12:15:43 7/18/2025 12:20:00 UPDATE

Resultat:

CompanyId Datapunkt __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULL 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 NULL
A XFv2 7/18/2025 12:15:43 NULL 7/18/2025 12:20:00 NULL

XFv1 ansågs giltig från 10:01:00 utan känt slut, och systemet höll den tron från 10:05:00 till 12:20:00. XFv1 är nu känt för att vara giltigt endast fram till 12:15:43, en korrigerad historik som gäller från systemtid 12:20:00 utan känt slut. XFv2 är giltigt från och med 12:15:43 utan kända slut och lärdes vid systemtid 12:20:00.

Ändring 3: Uppdatering i fel ordning

En uppdatering som kommer i fel ordning anländer och anger att Företag A faktiskt uppdaterades 2025-07-18 12:05:00 (verksamhetstid), men den tas inte in förrän 12:25:00 (systemtid). När en uppdatering anländer vid en senare systemtidpunkt men med en tidigare verksamhetstid korrigerar systemet den historiska verksamhetstiden och bevarar både det som systemet antog före den oordnade uppdateringen och den korrigerade historiken.

Input:

CompanyId Datapunkt Ordningsföljd Systemsekvensering Operation
A XFv3 7/18/2025 12:05:00 7/18/2025 12:25:00 UPDATE

Resultat:

CompanyId Datapunkt __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULL 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 NULL
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 NULL
A XFv2 7/18/2025 12:15:43 NULL 7/18/2025 12:20:00 NULL

XFv1 ansågs giltig från 10:01:00 till 12:15:43, och den tron är nu giltig i systemtid fram till 12:25:00. Den nya uppdateringen korrigerar XFv1:s affärsgiltighet så att den slutar kl. 12:05:00, och den korrigerade historiken gäller från systemtid 12:25:00. XFv3 är nu känt för att vara giltigt från 12:05:00 till 12:15:43, en tro som är giltig i systemtid från 12:25:00 utan känt slut.

Ändring 4: Ta bort

Företag A tas bort 2025-07-18 12:30:00 och systemet förbrukar händelsen kl. 12:30:00. Eftersom en borttagning innebär att entiteten upphör att existera i verksamheten skapar systemet ingen ersättningsrad. XFv2 visas i två rader, vilket bevarar en fullständig spårningslogg för både när företaget upphörde att existera och när systemet fick reda på borttagningen.

Input:

CompanyId Datapunkt Ordningsföljd Systemsekvensering Operation
A XFv2 7/18/2025 12:30:00 7/18/2025 12:30:00 DELETE

Resultat:

CompanyId Datapunkt __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULL 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 NULL
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 NULL
A XFv2 7/18/2025 12:15:43 NULL 7/18/2025 12:20:00 7/18/2025 12:30:00
A XFv2 7/18/2025 12:15:43 7/18/2025 12:30:00 7/18/2025 12:30:00 NULL

XFv2 var giltigt från 12:15:43 utan känt slut, och systemet höll den tron från 12:20:00 till 12:30:00. När borttagningen har registrerats är XFv2 känd för att vara giltig endast fram till 12:30:00, med en korrigerad historik med verkan från systemtid 12:30:00.

Vilka dataobjekt används för CDC-bearbetning i en pipeline?

När du deklarerar måltabellen i Hive-metaarkivet skapas två datastrukturer:

  • En vy med det namn som tilldelats måltabellen.
  • En intern stödtabell som används av pipelinen för att hantera CDC-bearbetning. Den här tabellen namnges genom att lägga till __apply_changes_storage_ framför måltabellens namn.

Om du till exempel deklarerar en måltabell med namnet dp_cdc_targetvisas en vy med namnet dp_cdc_target och en tabell med namnet __apply_changes_storage_dp_cdc_target i metaarkivet. Öppna vyn för att få tillgång till de bearbetade data. Ändra inte bakgrundstabellen direkt.

Anmärkning

Dessa datastrukturer gäller endast för AUTO CDC bearbetning, inte AUTO CDC FROM SNAPSHOT bearbetning. De gäller även endast hive-metaarkiv, inte Unity Catalog.