Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Important
Tato funkce je v beta verzi.
Tok REPLACE USING udržuje cílovou tabulku synchronizovanou se streamovacím zdrojem: nahradí všechny řádky, které odpovídají zadaným sloupcům klíče, a všechna ostatní data ponechá beze změny.
Sloupec řadí SEQUENCE BY aktualizace tak, aby výsledek byl správný i tehdy, když aktualizace přicházejí mimo pořadí. Pro každý klíč rozhoduje nejvyšší hodnota sekvence a řádek s nižší hodnotou sekvence nikdy nepřepíše řádek s vyšší hodnotou sekvence, který už je v cílové tabulce. Řádky, které sdílejí stejnou klávesu a stejnou sekvenci, jsou přidávány místo nahrazování.
Jak funguje REPLACE USING
Uvažujme tabulku událostí, která obsahuje klikání a konverzní události pro dvě oblasti, sekvenované podle seq:
| region_id | device_type | event_type | seq |
|---|---|---|---|
| 1 | iOS | click | 1 |
| 1 | Android | konverze | 1 |
| 2 | iOS | click | 1 |
| 2 | desktop | click | 1 |
Tok REPLACE USING (region_id) SEQUENCE BY seq přijímá tyto aktualizace pro regiony 1 a 3. Region 2 nemá žádné aktualizace:
| region_id | device_type | event_type | seq |
|---|---|---|---|
| 1 | iOS | click | 2 |
| 1 | Android | konverze | 2 |
| 1 | desktop | click | 2 |
| 3 | iOS | click | 1 |
| 3 | desktop | click | 2 |
Cílem je:
| region_id | device_type | event_type | seq | Outcome |
|---|---|---|---|---|
| 1 | iOS | click | 2 | Nahrazeno, protože sekv. 2 je větší než seq 1 |
| 1 | Android | konverze | 2 | Nahrazeno, protože sekv. 2 je větší než seq 1 |
| 1 | desktop | click | 2 | Nahrazeno, protože sekv. 2 je větší než seq 1 |
| 2 | iOS | click | 1 | Nedotčený, protože klíč v této aktualizaci není přítomen |
| 2 | desktop | click | 1 | Nedotčený, protože klíč v této aktualizaci není přítomen |
| 3 | desktop | click | 2 | Přidáno. Řádek sekv. 1 pro oblast 3 není přidán, protože se aplikuje pouze nejvyšší sekvence pro klíč. |
Požadavky
Toky REPLACE USING mají následující požadavky:
- Toky REPLACE USING běží v Databricks Runtime 18.2 a vyšším na klasických nebo bezserverových výpočetních prostředcích. Databricks doporučuje Unity Catalog.
- Zdroj musí být streamovací zdroj. REPLACE USING odmítá neproudový zdroj.
- Musíte zadat alespoň jeden klíčový sloupec a přesně jeden
SEQUENCE BYsloupec.
Kdy použít VÝMĚNU POMOCÍ toků
Kanály Lakeflow nabízejí tři toky, které přepisují stávající řádky. Vyberte podle toho, jak váš zdroj vypadá a jak identifikuje řádky k nahrazení:
- Použijte REPLACE USING, když je zdrojem řada částečných snímků identifikovaných podle sloupce. REPLACE USING přepíše pouze data, pro která v příchozích datech existuje odpovídající záznam, a všechna ostatní data ponechá beze změny. Nevyžaduje primární klíč.
- Používejte AUTO CDC, pokud je vaším zdrojem feed pro zachycování změn dat (CDC) s explicitními operacemi vkládání, aktualizace a mazání, nebo pokud potřebujete historii typu 2 pomalu se měnící dimenze (SCD). AUTO CDC také vyžaduje skutečný primární klíč. Viz rozhraní API AUTO CDC: Zjednodušte zachytávání změn dat pomocí pipelin.
- Použijte REPLACE WHERE, když je zdroj snapshot a chcete přepočítat a přepsat rozsah cílové tabulky určený predikátem, například za posledních 7 dní, v rámci dávkové operace. Nevyžaduje primární klíč. Viz Dávkové zpracování pomocí toků REPLACE WHERE.
Vytvořte tok pomocí REPLACE USING
Definujte procesy REPLACE USING buď v SQL, nebo v Pythonu.
SQL
Použijte klauzuli FLOW REPLACE USING přímo v řádku s 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);
Alternativně použijte syntaxi dlouhého formátu 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 vyžaduje se v SQL. Přiřazuje sloupce podle názvu spíše než podle pozice.
Python
Označte tabulku a tok společně s @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")
Případně zacílit na existující streamovací tabulku pomocí @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 je seznam klíčových sloupců.
sequence_by je název sloupce nebo Column výraz a je vyžadován vždy, když replace_using je nastaven.
Řazení a data v nesprávném pořadí
Sloupec SEQUENCE BY činí výsledek nezávislým na pořadí, v jakém aktualizace přicházejí. Řádek se aplikuje na klíč pouze tehdy, pokud je jeho posloupnost větší než sekvence již uložená pro daný klíč, takže pozdní nebo opakovaně přehraný řádek, který je starší než aktuální hodnota, je ignorován. Klíče, které v aktualizaci nejsou, zůstávají nedotčené.
Dodržujte tyto postupy, aby se náhrada chovala předvídatelně:
| Postupy | Reason |
|---|---|
| Použijte posloupnost, která se pro každou verzi klíče striktně zvyšuje, například časové razítko, číslo verze nebo offset v protokolu. | Dva řádky se stejným klíčem a stejnou sekvencí jsou zachovány, což vede k duplicitním řádkům pro daný klíč. |
| Použijte nenulovou sekvenci. | Nulová sekvence může vést k nedefinovanému chování. |
Expectations
NAHRADIT POMOCÍ toků očekávání podpory.
warn a fail chovají se jako u jiných toků: warn pokračují v porušování řádků, zaznamenávají porušení a fail zastaví aktualizaci. Viz Spravujte kvalitu dat pomocí požadavků na datový potrubí.
Očekávání drop považuje řádek porušující toto očekávání za takový, jako by jej zdroj nikdy nevytvořil. Vynechaný řádek nenahrazuje, nemazá ani nemění odpovídající klíče v cílové tabulce:
- Vyhazování nastává před deduplikací, takže flow uchovává nejnovější platnou verzi klíče.
- Pokud je každý příchozí řádek pro klíč vynechán, stávající řádky klíče zůstanou nedotčeny.
- Protože zahozený řádek nenastavuje žádnou dolní mez sekvence, pozdější platná aktualizace se přesto použije, i když má nižší sekvenční číslo než zahozený řádek.
Omezení
Toky „Nahradit pomocí“ mají následující omezení:
- REPLACE USING podporuje pouze jeden tok pro každou cílovou tabulku. Kombinování REPLACE USING s jiným typem toku na stejném cíli není podporováno.
- Cílová tabulka musí být vytvořena v rámci pipeline.
- Zdroj musí být streamovací.
- Musíte zadat alespoň jeden klíčový sloupec a jeden
SEQUENCE BYsloupec. Klíčové sloupce nelze opakovat a typ každého klíčového sloupce musí být tříditelný. Atomické typy, jako jsou celá čísla, řetězce a data, mohou být klíče, zatímcoMAPaVARIANTnemohou. - Pro samostatné streamovací tabulky viz Apply partial snapshot replacement with REPLACE USING flows, kde najdete rozdíly v syntaxi.
Examples
Následující příklady jsou z samples.wanderbricks.booking_updates, ukázkové tabulky změn stavu rezervací, která je dostupná ve všech pracovních prostorech s podporou Unity Catalog. Každá rezervace se objeví jednou za každou změnu, takže booking_id se opakuje s novým booking_update_id. Viz datovou sadu Wanderbricks.
Příklad 1: Uchovejte nejnovější záznam pro každý klíč
Uchovejte pouze aktuální stav každé rezervace. Tok je klíčován podle booking_id a řazen podle booking_update_id, takže nejnovější aktualizace rezervace nahradí ty předchozí. Používejte místo toho AUTO CDC, pokud je zdrojem feed změn s explicitními operacemi vkládání, aktualizace a mazání.
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")
Tento příklad se řadí podle booking_update_id, nikoli podle časového razítka updated_at, protože více aktualizací téže rezervace může sdílet stejné časové razítko. Řádky se stejnou hodnotou pořadí jsou přidány, nikoli nahrazeny, takže pro tyto rezervace zůstane více než jeden řádek.
Příklad 2: Zadejte více než jeden sloupec
Když je záznam identifikován kombinací sloupců, uveďte je všechny v .REPLACE USING Zde je každá rezervace identifikována jako (property_id, booking_id), takže tok zachovává aktuální stav každé rezervace na nemovitost. Pokud může být sloupec klíče null, REPLACEMENT USING porovnává null s null místo přeskakování řádku.
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")
Příklad 3: Vyřaďte neplatné záznamy s očekáváním
Přidejte pravidlo, které zabrání tomu, aby se vadné řádky dostaly do cílového umístění. Vynechaný řádek se považuje, jako by ho zdroj nikdy nevytvořil: nenahrazuje ani nesmaže odpovídající klíč a tok se vrací k nejposlednímu platnému řádku pro daný klíč. Tento flow vypouští aktualizace, které nemají pozitivní 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")