Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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 BYoszlopot 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 BYegy 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ígMAPésVARIANTnem 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")