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.
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 BYkolumn.
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 BYkolumn. 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, medanMAPochVARIANTinte 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")