Ersättning av partiell ögonblicksbild med REPLACE USING-flöden

Important

Den här funktionen finns i Beta.

Ett REPLACE USING-flöde håller en måltabell synkroniserad med en strömmande källa: det ersätter alla rader som matchar de angivna nyckelkolumnerna och lämnar övriga data oförändrade.

En SEQUENCE BY kolumn ordnar uppdateringarna så att resultatet är korrekt även när uppdateringar kommer in i fel ordning. För varje nyckel är det den högsta sekvensen som gäller, och en rad med lägre sekvensnummer skriver aldrig över en rad med högre sekvensnummer som redan finns i målet. Rader som har samma nyckel och samma sekvens läggs till i stället för att ersättas.

Hur REPLACE USING fungerar

Betrakta en händelsetabell som håller klick- och konverteringshändelser för två regioner, sekvenserade av seq:

region_id enhetstyp event_type seq
1 iOS klicka 1
1 Android omvandling 1
2 iOS klicka 1
2 Skrivbordet klicka 1

Ett REPLACE USING (region_id) SEQUENCE BY seq flöde får dessa uppdateringar för regionerna 1 och 3. Region 2 har inga uppdateringar:

region_id enhetstyp händelsetyp seq
1 iOS klicka 2
1 Android omvandling 2
1 Skrivbordet klicka 2
3 iOS klicka 1
3 Skrivbordet klicka 2

Målet blir:

region_id enhetstyp event_type seq Outcome
1 iOS klicka 2 Ersatt, eftersom seq 2 är större än seq 1
1 Android omvandling 2 Ersatt, eftersom seq 2 är större än seq 1
1 Skrivbordet klicka 2 Ersatt, eftersom seq 2 är större än seq 1
2 iOS klicka 1 Oförändrad, eftersom nyckeln inte finns i denna uppdatering
2 Skrivbordet klicka 1 Oförändrad, eftersom nyckeln inte finns i den här uppdateringen
3 Skrivbordet klicka 2 Tillagd. Seq 1-raden för region 3 läggs inte till, eftersom endast den högsta sekvensen för en nyckel tillämpas.

Requirements

REPLACE USING-flöden har följande krav:

  • Ersätt med hjälp av flöden som körs i Databricks Runtime 18.2 och senare, på klassiska eller serverlösa beräkningsresurser. Databricks rekommenderar Unity Catalog.
  • Källan måste vara en strömmande källa. REPLACE USING avvisar en icke-strömmande källa.
  • Du måste ange minst en nyckelkolumn och exakt en SEQUENCE BY kolumn.

När används REPLACE USING-flöden

Lakeflow pipelines erbjuder tre flöden som ersätter befintliga rader. Välj baserat på hur din källa ser ut och hur den identifierar vilka rader som ska ersättas:

  • Använd REPLACE USING när din källa är en serie delvisa ögonblicksbilder med kolumn som nyckel. REPLACE USING skriver bara över den data som har en matchning i den inkommande datan, och lämnar all annan data orörd. Det kräver ingen primärnyckel.
  • Använd AUTO CDC när din källa är ett CDC-flöde (change data capture) med explicita infogningar, uppdateringar och borttagningar, eller när du behöver historik för långsamt föränderliga dimensioner av typ 2 (SCD). AUTO CDC kräver också en riktig primärnyckel. Se API:er för AUTOMATISK CDC: Förenkla insamling av ändringsdata med pipelines.
  • Använd REPLACE WHERE när din källa är en snapshot och du vill beräkna om och skriva över ett intervall i måltabellen som valts av ett predikat, till exempel de senaste 7 dagarna, som en batchoperation. Det kräver ingen primärnyckel. Se batchbearbetning med REPLACE WHERE-flöden.

Skapa ett REPLACE USING-flöde

Definiera REPLACE USING-flöden med antingen SQL eller Python.

SQL

Använd FLOW REPLACE USING satsen i samma rad som CREATE STREAMING TABLE:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Du kan också använda långformssyntaxen CREATE FLOW :

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME krävs i SQL. Den matchar kolumner efter namn i stället för position.

Python

Deklarera tabellen och flödet tillsammans med @dp.table:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Alternativt kan du rikta dig mot en befintlig streamingtabell med @dp.replace_flow:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using är en lista över nyckelkolumner. sequence_by är ett kolumnnamn eller ett Column uttryck, och krävs när replace_using som helst är satt.

Sekvensering och data i fel ordning

Kolumnen SEQUENCE BY gör resultatet oberoende av i vilken ordning uppdateringarna anländer. En rad tillämpas på en nyckel endast om dess sekvens är större än den sekvens som redan lagrats för den nyckeln, så en sen eller omspelad rad som är äldre än det aktuella värdet ignoreras. Nycklar som inte finns i en uppdatering lämnas orörda.

Följ dessa metoder så att ersättningen beter sig förutsägbart:

Practice Förnuft
Använd en sekvens som strikt ökar per nyckelversion, såsom tidsstämpel, versionsnummer eller loggoffset. Två rader med samma nyckel och samma sekvens sparas båda, vilket resulterar i dubbla rader för den nyckeln.
Använd en icke-null sekvens. En nollsekvens kan leda till odefinierat beteende.

Expectations

ERSÄTT MED flöden som stödjer förväntningarna. warn fail och beter sig som de gör i andra flöden: warn fortsätter att bryta mot rader och registrerar överträdelsen, och fail stoppar uppdateringen. Se avsnittet Hantera datakvalitet med pipeline-förväntningar.

En drop förväntan behandlar en rad som bryter mot förväntningen som om källan aldrig hade genererat den. Den borttagna raden ersätter, tar inte bort eller ändrar matchande nycklar i måltabellen:

  • Borttagning sker före deduplicering, så flödet behåller den senaste giltiga versionen för den nyckeln.
  • Om varje inkommande rad för en nyckel tas bort, lämnas nyckelns befintliga rader orörda.
  • Eftersom en bortsläppt rad inte anger någon sekvensgolv, landar en senare giltig uppdatering ändå även om dess sekvens är lägre än den borttagna radens.

Limitations

REPLACE USING-flödena har följande begränsningar:

  • REPLACE USING stöder ett enda flöde per måltabell. Det går inte att kombinera REPLACE USING med en annan flödestyp med samma mål.
  • Måltabellen måste skapas inom pipelinen.
  • Källan måste vara en strömmande källa.
  • Du måste specificera minst en nyckelkolumn och en SEQUENCE BY kolumn. Nyckelkolumner kan inte upprepas, och varje nyckelkolumns typ måste vara sorterbar. Atomära typer, såsom heltal, strängar och datum, kan vara nycklar, medan MAP och VARIANT inte kan.
  • För fristående streamingtabeller, se Tillämpa partiell ersättning av ögonblicksbild med REPLACE USING-flöden för information om skillnader i syntaxen.

Examples

Följande exempel hämtas från samples.wanderbricks.booking_updates, en exempeltabell med ändringar i bokningsstatus som finns tillgänglig i varje arbetsyta med Unity Catalog aktiverat. Varje bokning visas en gång per förändring, så booking_id upprepas med en ny booking_update_id. Se Wanderbricks dataset.

Exempel 1: Spara den senaste posten för varje nyckel

Behåll endast det aktuella läget för varje bokning. Flödet använder booking_id som nyckel och ordnas i sekvens efter booking_update_id, så den senaste uppdateringen för en bokning ersätter tidigare uppdateringar. Använd istället AUTO CDC när din källa är ett ändringsflöde med explicita insättnings-, uppdaterings- och borttagningsoperationer.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Detta exempel sekvenserar efter booking_update_id istället för updated_at tidsstämpel eftersom flera uppdateringar av samma bokning kan dela en tidsstämpel. Rader som har samma sekvensnummer läggs till i stället för att ersättas, vilket innebär att det blir mer än en rad för dessa bokningar.

Exempel 2: Nyckel på mer än en kolumn

När en post identifieras av en kombination av kolumner, lista dem alla i REPLACE USING. Här identifieras varje bokning med (property_id, booking_id), så flödet behåller det aktuella tillståndet för varje bokning per fastighet. Om en nyckelkolumn kan vara null matchar REPLACE USING null med null snarare än att hoppa över raden.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Exempel 3: Ta bort ogiltiga poster med en förväntning

Lägg till en förväntan om att hålla dåliga rader borta från målet. En borttagen rad behandlas som om källan aldrig producerat den: den ersätter eller tar inte bort den matchande nyckeln, och flödet faller tillbaka till den senaste giltiga raden för den nyckeln. Detta flöde ignorerar uppdateringar som inte har ett positivt total_amount.

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")