Részleges snapshot helyettesítés REPLACE USING flow-okkal

Important

Ez a funkció bétaverzióban érhető el.

A REPLACE USING flow szinkronban tartja a céltáblát egy streaming forrással: lecseréli az összes sort, amely megfelel a megadott kulcsoszlopnak, és minden egyéb adatot változatlanul hagy.

Egy SEQUENCE BY oszlop rendeli a frissítéseket, így az eredmény helyes akkor is, ha a frissítések nem sorrendben érkeznek. Minden kulcsnál a legmagasabb sorozat nyer, és egy alacsonyabb sorrendű sor soha nem írja felül a célban már meglévő magasabb sort. Azokat a sorokat, amelyek ugyanazt a hangnemet és sorozatot használják, nem cserélik ki, hanem hozzáadják.

Hogyan működik a REPLACE USING

Vegyünk egy eseménytáblázatot, amely két régióra kattintás és konverziós eseményeket tartalmaz, a következőképpen seqsorrendben:

region_id eszköztípus event_type seq
1 iOS click 1
1 Android átalakítás 1
2 iOS click 1
2 asztal click 1

Egy REPLACE USING (region_id) SEQUENCE BY seq áramlás kapja ezeket a frissítéseket az 1. és 3. régióra. A 2-es régióban nincs frissítés:

region_id eszköztípus eseménytípus seq
1 iOS click 2
1 Android átalakítás 2
1 asztal click 2
3 iOS click 1
3 asztal click 2

A célpont a következő lesz:

region_id eszköztípus event_type seq Outcome
1 iOS click 2 Lecserélve, mert a 2. szekció nagyobb, mint az 1.
1 Android átalakítás 2 Lecserélve, mert a 2. szekció nagyobb, mint az 1.
1 asztal click 2 Kicserélve, mert a seq 2 nagyobb, mint a seq 1.
2 iOS click 1 Érintetlen, mert a kulcs nincs jelen ebben a frissítésben
2 asztal click 1 Érintetlen, mert a kulcs nincs jelen ebben a frissítésben
3 asztal click 2 Hozzáadva. A 3-as régióhoz tartozó 1-es szekvenciájú sor nem kerül hozzáadásra, mert egy kulcshoz csak a legmagasabb szekvencia kerül alkalmazásra.

Követelmények

A HELYETTESÍTÉS HASZNÁLATA folyamatokra az alábbi követelmények vonatkoznak:

  • A REPLACE USING műveletet használó folyamatok a Databricks Runtime 18.2-es vagy újabb verzióján, klasszikus vagy serverless számítási környezetben futnak. A Databricks a Unity Catalog-t ajánlja.
  • A forrásnak streamelési forrásnak kell lennie. A REPLACE USING elutasít egy nem streaming forrást.
  • Legalább egy kulcsoszlopot és pontosan egy SEQUENCE BY oszlopot kell megadnod.

Mikor kell használni a REPLACE USING folyamatokat

A Lakeflow-folyamatok három adatfolyamot kínálnak, amelyek a meglévő sorokat írják felül. Válassz az alapján, hogy néz ki a forrásod, és hogyan azonosítja a sorokat, amelyeket a helyettesíteni kell:

  • Használd a REPLACE USING módot, ha a forrásod részleges snapshotok sorozata, amelyeket oszloponként kulcsolnak. A REPLACE USING csak azokat az adatokat írja felül, amelyek egyeznek a bejövő adatokban, így minden más adat érintetlen marad. Nem igényel elsődleges kulcsot.
  • Használd az AUTO CDC-t, ha a forrásod egy változás-adat-capture (CDC) adatfolyam, amely explicit beillesztés, frissítés és törlés műveleteket tartalmaz, vagy lassan kell változtatni a dimenziós (SCD) 2-es típusú előzményeket. Az AUTO CDC is megköveteli a valódi elsődleges kulcsot. Lásd az AUTO CDC API-k: Egyszerűsítse a változáskövető adatrögzítést a csővezetékekkel.
  • Használd a REPLACE WHERE parancsot, amikor a forrás egy pillanatkép, és egy predikátum alapján kiválasztott tartományt — például az elmúlt 7 napot — szeretnél újraszámítani és felülírni a céltáblában, kötegelt műveletként. Nem igényel elsődleges kulcsot. Lásd: Kötegelt feldolgozás REPLACE WHERE folyamatokkal.

Hozz létre egy REPLACE USING folyamatot

Definiáld a REPLACE USING folyamatokat SQL-ben vagy Python-ban.

SQL

Használja a(z) FLOW REPLACE USING záradékot egy sorban a(z) CREATE STREAMING TABLE elemmel:

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

Másik lehetőségként használja a hosszú formátumú szintaxist 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);

Megjegyzés:

BY NAME kötelező az SQL-ben. Az oszlopokat név alapján párosítja, nem pedig pozíció alapján.

Python

Deklaráljuk a táblát és az áramlást együtt:@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")

Alternatívaként célozz meg egy meglévő streaming táblát a következőkkel @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 a kulcsoszlopok listája. sequence_by oszlopnév vagy Column kifejezés, és minden beállításkor replace_using kötelező.

Szekvenálás és nem sorrendben érkező adatok

Az SEQUENCE BY oszlop független az eredményt attól, hogy milyen sorrendben érkeznek a frissítések. Egy sor csak akkor kerül egy kulcsra, ha annak sorrendje nagyobb, mint az adott kulcshoz már tárolt sorozat, így egy későn vagy újrajátszott sort, amely régebbi, mint a jelenlegi érték, figyelmen kívül kerül. A frissítésekben nem található kulcsok érintetlenek maradnak.

Kövesd ezeket a gyakorlatokat, hogy a csere kivárhatóan viselkedjen:

Practice Indok
Használj minden kulcsverzióhoz olyan sorozatot, amely szigorúan növekszik, például időbélyeget, verziószámot vagy naplóeltolást. Két, azonos kulccsal és azonos sorszámmal rendelkező sor mindkettő megmarad, ami duplikált sorokat eredményez ennél a kulcsnál.
Használj egy nem null sorozatot. Egy null szekvencia meghatározatlan viselkedéshez vezethet.

Expectations

REPLACE USING folyamatok megfelelnek az elvárásoknak. warn és fail úgy viselkednek, ahogy más folyamatoknál: warn folyamatosan megsértik a sorokat, rögzítik a szabálysértést, és fail megállítja a frissítést. Lásd Az adatminőség kezelése folyamatelvárásokkal.

Egy drop elvárás úgy kezeli a megsértő sort, mintha a forrás soha nem hozta volna létre. Az elhagyott sor nem helyettesíti, törli vagy módosítja a megfelelő kulcsokat a céltáblában:

  • Az eldobás a deduplikáció előtt történik, így a folyamat a kulcshoz tartozó legfrissebb érvényes verziót tartja meg.
  • Ha egy adott kulcshoz tartozó összes bejövő sor eldobásra kerül, a kulcshoz tartozó meglévő sorok érintetlenek maradnak.
  • Mivel egy kiesett sor nem állít be szekvencia szintet, egy későbbi érvényes frissítés akkor is érkezik, ha annak sorozata alacsonyabb, mint az elhagyott soré.

Limitations

A CSERE HASZNÁLATÁVAL folyamatoknak az alábbi korlátaik vannak:

  • A REPLACE USING céltáblánként csak egy folyamatot támogat. A REPLACE USING kombinálása egy másik folyamattípussal ugyanazon a célponton nem támogatott.
  • A céltáblát a folyamaton belül kell létrehozni.
  • A forrásnak streamelési forrásnak kell lennie.
  • Legalább egy kulcsoszlopot és SEQUENCE BY egy oszlopot kell megadnod. A kulcsoszlopokat nem lehet ismételni, és minden kulcsoszlop típusának rendezhetőnek kell lennie. Az atomtípusok, mint az egész számok, láncsorok és dátumok, lehetnek kulcsok, míg MAP és VARIANT nem lehetnek.
  • Az önálló streaming táblák esetében a szintaktikai különbségeket illetően lásd: Apply partial snapshot replacement with REPLACE USING flows.

Examples

Az alábbi példák a foglalási állapotváltozásokat tartalmazó samples.wanderbricks.booking_updates mintatáblából olvasnak, amely minden Unity Catalogdal rendelkező munkaterületen elérhető. Minden foglalás minden módosításnál egyszer jelenik meg, így booking_id újra megjelenik egy új booking_update_id mellett. Lásd Wanderbricks adathalmazát.

1. példa: Tartsd a legfrissebb rekordot minden billentyűhöz

Csak az aktuális foglalási állapotot tartsd meg. A folyamat a booking_id alapján kulcsol, és a booking_update_id alapján sorrendbe állítja az elemeket, így egy foglalás legfrissebb frissítése felülírja a korábbiakat. Használd helyette az AUTO CDC-t, ha a forrásod egy változáscsatorna, amely explicit beszúrási, frissítési és törlési műveleteket tartalmaz.

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

Ez a példa a booking_update_id alapján állítja sorrendbe az elemeket, nem pedig a updated_at időbélyeg alapján, mert ugyanahhoz a foglaláshoz tartozó több frissítés is kaphat azonos időbélyeget. Az azonos sorszámú sorokat hozzáfűzik ahelyett, hogy lecserélnék őket, így ezekhez a foglalásokhoz egynél több sor maradna.

Példa 2: Több oszlopra kell kattintani

Ha egy rekordot oszlopok kombinációjával azonosítunk, sorold fel őket a .REPLACE USING Itt minden foglalást az (property_id, booking_id)alapján azonosítanak, így a folyamat megőrzi minden foglalás aktuális állapotát minden ingatlanonként. Ha egy kulcsoszlop értéke NULL lehet, a REPLACE USING a NULL értéket a NULL értékhez illeszti, ahelyett hogy kihagyná a sort.

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

3. példa: Érvénytelen rekordok elhagyása egy elvárással

Tedd hozzá az elvárást, hogy a rossz sorok távol maradjanak a célponttól. Az elhagyott sort úgy kezelik, mintha a forrás soha nem hozta volna létre: nem cseréli vagy törli a megfelelő kulcsot, és a folyamat visszakerül a legfrissebb érvényes sorra az adott kulcshoz. Ez a folyamat olyan frissítéseket dob el, amelyeknek nincs pozitív total_amounteredményük.

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