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.
AUTO CDC ... INTO Az utasítás használatával hozzon létre egy folyamatot, amely a Lakeflow-folyamatokat használja az adatrögzítési (CDC) funkciók módosításához. Ez az utasítás beolvassa a CDC-forrás módosításait, és alkalmazza őket egy streamelési célra.
- A CDC-vel kapcsolatos további információkért lásd: Adatrögzítés és pillanatképek módosítása.
- A
AUTO CDChasználatáról további információkat a „AUTO CDC API-k: Az adatrögzítés egyszerűsítése folyamatokkal” című témakörben talál. - További részletekért
CREATE FLOWlásd a CREATE FLOW (pipeline-ek) című témakört.
Szemantika
CREATE OR REFRESH STREAMING TABLE table_name;
CREATE FLOW flow_name AS AUTO CDC [ONCE] INTO table_name
FROM source
KEYS (keys)
[IGNORE NULL UPDATES [ON {columnList | * EXCEPT (exceptColumnList)}]]
[APPLY AS DELETE WHEN condition]
[APPLY AS TRUNCATE WHEN condition]
SEQUENCE BY orderByColumn
[SYSTEM SEQUENCE BY systemOrderByColumn]
[COLUMNS {columnList | * EXCEPT (exceptColumnList)}]
[STORED AS {SCD TYPE 1 | SCD TYPE 2 | BITEMPORAL}]
[TRACK HISTORY ON {columnList | * EXCEPT (exceptColumnList)}]
[COLUMNS TO UPDATE columnName]
A cél adatminőségi korlátozásait ugyanazzal CONSTRAINT a záradékkal határozhatja meg, mint a többi folyamat lekérdezését. Lásd: Adatminőség kezelése folyamatelvárásokkal.
Az INSERT és UPDATE események alapértelmezett viselkedése az, hogy a forrásból származó CDC-eseményeket frissíti vagy beszúrja: frissíti a céltábla azon sorait, amelyek megfelelnek a megadott kulcs(ok)nak, vagy új sort szúr be, ha egyező rekord nem található a céltáblában.
DELETE események kezelése a APPLY AS DELETE WHEN feltétellel meghatározható.
Fontos
A módosítások alkalmazásához deklarálnia kell egy célstreamelési táblát. Igény szerint megadhatja a céltábla sémáját. A 2. típusú SCD-táblák esetében a céltábla sémájának megadásakor a __START_AT mezővel azonos adattípusú __END_AT és sequence_by oszlopokat is tartalmaznia kell.
Lásd az AUTO CDC API-k: Egyszerűsítse a változáskövető adatrögzítést a csővezetékekkel.
Paraméterek
ONCEA beállítás
ONCEazt jelenti, hogy ez egyszeri beszúrást vagy visszatöltést végez a céltáblába. Nem fut újra, ha a benne lévő folyamat frissül, kivéve a teljes frissítést.Ez a záradék nem kötelező.
flow_nameA létrehozandó folyamat neve.
sourceAz adatok forrása. A forrásnak streamelési forrásnak kell lennie. A STREAM kulcsszóval stream-szemantikát használhat a forrásból való olvasáshoz. Ha az olvasás egy meglévő rekord módosítását vagy törlését tapasztalja, hibaüzenet jelenik meg. A legbiztonságosabb, ha statikus vagy csak hozzáfűző forrásokból olvas. A módosítási véglegesítéseket tartalmazó adatok betöltéséhez használhatja a Pythont és a
skipChangeCommitshibák kezelésére szolgáló lehetőséget.További információért az adatfolyamokról, tekintse meg a Folyamatokkal történő adatátalakítást.
KEYSAzok az oszlopok vagy oszlopok kombinációja, amelyek egyedileg azonosítják a forrásadatok sorait. Az oszlopokban szereplő értékek segítségével azonosíthatja, hogy mely CDC-események vonatkoznak a céltábla adott rekordjaira.
Az oszlopok kombinációjának meghatározásához használjon vesszővel tagolt oszloplistát.
Ez a záradék kötelező.
IGNORE NULL UPDATESLehetővé teszi a céloszlopok egy részhalmazát tartalmazó frissítések betöltését. Ha egy CDC-esemény megfelel egy meglévő sornak, és
IGNORE NULL UPDATESmeg van adva, az értékekkel rendelkezőnulloszlopok megőrzik a célban meglévő értékeiket. Ez anullértékkel rendelkező beágyazott oszlopokra is vonatkozik.Részleges frissítések esetén adjon hozzá egy záradékot
ON, amely szabályozza, hogy mely oszlopok figyelmen kívül hagyjáknullaz értékeket:-
IGNORE NULL UPDATES ON columnList: csak a felsorolt oszlopok őrzik meg a meglévő értékeket, ha a bejövő érték .nullMinden más oszlop explicitnullértékeket alkalmaz. -
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList): a felsoroltak kivételével minden oszlop megőrzi a meglévő értékeit, ha a bejövő érték .nullA felsorolt oszlopok explicitnullértékeket alkalmaznak.
Ez a záradék nem kötelező.
Az alapértelmezett az, hogy a meglévő oszlopokat a
nullértékekkel írjuk felül.-
APPLY AS DELETE WHENMegadja, hogy a CDC-eseményeket mikor kell
DELETE-ként, és nem upsertként kezelni.A 2. típusú SCD-források esetében a rendezetlen adatok kezelésére a rendszer ideiglenesen megőrzi a törölt sort a háttérben lévő Delta-táblában, és létrehoz egy nézetet a metaadattárban, amely kiszűri ezeket a sírköveket. A megőrzési időköz a
pipelines.cdc.tombstoneGCThresholdInSecondshasználatával konfigurálható.Ez a záradék nem kötelező.
APPLY AS TRUNCATE WHENMegadja, hogy a CDC-események mikor legyenek teljes táblaként
TRUNCATEkezelve. Mivel ez a záradék a céltábla teljes csonkját aktiválja, csak a funkciót igénylő konkrét használati esetekhez használható.A
APPLY AS TRUNCATE WHENzáradék csak az 1. SCD-típus esetében támogatott. A 2. SCD-típus nem támogatja a csonkolási műveletet.Ez a záradék nem kötelező.
SEQUENCE BYAz oszlop neve, amely a CDC-események logikai sorrendjét adja meg a forrásadatokban. A pipeline feldolgozás ezzel a szekvenálással kezeli a nem sorrendben érkező változási eseményeket.
Ha több oszlopra van szükség a szekvenáláshoz, használjon egy
STRUCTkifejezést: először az első struct mező, majd a második mező alapján rendezi, ha döntetlen van, és így tovább.A megadott oszlopoknak rendezhető adattípusoknak kell lenniük.
Ez a záradék kötelező.
SYSTEM SEQUENCE BYFontos
A Bitemporal AUTO CDC bétaverzióban érhető el.
Az oszlop neve, amely azt a rendszeridőt adja meg, amikor az egyes CDC-események ismertek a rendszer számára. Az üzleti idő (
STORED AS BITEMPORAL) és a rendszeridő változásainak nyomon követéséreSEQUENCE BYszolgál. Lásd: Bitemporal AUTO CDC.A megadott oszlopoknak rendezhető adattípusoknak kell lenniük.
Ez a záradék nem kötelező, és csak bitemporális táblákra vonatkozik.
COLUMNSA céltáblában szerepeltetni kívánt oszlopok egy részhalmazát adja meg. A következő lehetőségek közül választhat:
- Adja meg a belefoglalandó oszlopok teljes listáját:
COLUMNS (userId, name, city). - Adja meg a kizárandó oszlopok listáját:
COLUMNS * EXCEPT (operation, sequenceNum)
Ez a záradék nem kötelező.
Az alapértelmezett érték az, hogy az összes oszlopot belefoglalja a céltáblába, ha a
COLUMNSzáradék nincs megadva.- Adja meg a belefoglalandó oszlopok teljes listáját:
STORED ASA rekordok tárolása SCD 1- vagy 2-es SCD-típusként vagy bitemporálisként.
Úgy van beállítva, hogy
BITEMPORALnyomon kövesse az üzleti idő és a rendszeridő változásait. A Bitemporal megköveteliSYSTEM SEQUENCE BYés bétaverzióban van. Lásd: Bitemporal AUTO CDC.Ez a záradék nem kötelező.
Az alapértelmezett scd típus 1.
TRACK HISTORY ONA kimeneti oszlopok egy részhalmazát adja meg, amely előzményrekordokat hoz létre a megadott oszlopok módosításakor. A következő lehetőségek közül választhat:
- Adja meg a követendő oszlopok teljes listáját:
COLUMNS (userId, name, city). - Adja meg a nyomon követésből kizárandó oszlopok listáját:
COLUMNS * EXCEPT (operation, sequenceNum)
Ez a záradék nem kötelező. Az alapértelmezett beállítás az összes kimeneti oszlop előzményeinek nyomon követése, ha bármilyen változás történt, ami egyenértékű azokkal
TRACK HISTORY ON *.- Adja meg a követendő oszlopok teljes listáját:
COLUMNS TO UPDATEMegadja annak a forrásoszlopnak a nevét, amely az egyes változásrekordok esetében oszlopnév-sztringek tömbjeként (
array<string>) tárolja a frissíteni kívánt oszlopokat. A tömbben nem szereplő oszlopok megtartják a meglévő célértékeiket, míg a felsorolt oszlopok a forrásból vannak megírva, beleértve az explicitnullértékeket is.Ez a záradék részleges frissítésekhez használható, ha minden változásrekord más oszlopkészletet frissít, és explicit
nullértékeket kell alkalmaznia.Nem használható
COLUMNS TO UPDATEegyütt,IGNORE NULL UPDATESés bitemporális táblák esetében nem támogatott.Ez a záradék nem kötelező.
Példák
-- Create a streaming table, then use AUTO CDC to populate it:
CREATE OR REFRESH STREAMING TABLE target;
CREATE FLOW flow
AS AUTO CDC INTO
target
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2
TRACK HISTORY ON * EXCEPT (city);