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.
A Lakeflow-folyamatcsatornák leegyszerűsítik a változásadat-rögzítést (CDC) az AUTO CDC és AUTO CDC FROM SNAPSHOT API-kkal. Ezek az API-k automatizálják a lassan változó dimenziók (SCD) 1-es és 2-es típusának feldolgozását, akár CDC forrásokból, akár adatbázis-pillanatképekből kiindulva. Az AUTO CDC API támogatja a bitemporális nyomkövetést is, amely két idődimenzió (bétaverzió) változásait rögzíti. Az 1. és a 2. típusú SCD-ről az adatrögzítés és a pillanatképek módosítása című témakörben olvashat bővebben. További információ a bitemporális nyomkövetésről: Bitemporal AUTO CDC.
Megjegyzés:
Az AUTO CDC API-k lecserélik az APPLY CHANGES API-kat, és ugyanazzal a szintaxisukkal rendelkeznek. Az APPLY CHANGES API-k továbbra is elérhetők, de a Databricks az AUTO CDC API-k használatát javasolja helyettük.
A használt API a változásadatok forrásától függ:
-
AUTO CDC: Ezt akkor használja, ha a forrásadatbázisban engedélyezve van a CDC-hírcsatorna.AUTO CDCegy változásadatcsatorna (CDF) változásait dolgozza fel. Ez a folyamat SQL- és Python-felületeiben is támogatott. -
AUTO CDC FROM SNAPSHOT: Ezt akkor használja, ha a CDC nincs engedélyezve a forrásadatbázisban, és csak pillanatképek érhetők el. Ez az API összehasonlítja a pillanatképeket a módosítások meghatározásához, majd feldolgozza őket.
Mindkét API támogatja a táblák frissítését az 1. és a 2. típusú SCD használatával:
- Az SCD 1. típusának használatával közvetlenül frissítheti a rekordokat. A rendszer nem őrzi meg a frissített rekordok előzményeit.
- Az SCD 2. típusával megőrizheti a rekordok előzményeit, akár az összes frissítésen, akár egy adott oszlopkészlet frissítésén.
Csak AUTO CDC esetén használható a bitemporális tárolás is, amely kiterjeszti az SCD 2-es típusú előzménykezelését, hogy két idődimenzió mentén kövesse nyomon a változásokat: az üzleti idő és a rendszeridő szerint. A Bitemporal bétaverzióban van. Lásd Bitemporal AUTO CDC.
AUTO CDC A részleges frissítéseket is támogatja, ahol a változásrekord csak az oszlopok egy részét frissíti. Lásd: Részleges frissítések alkalmazása.
Az AUTO CDC API-kat nem támogatják az Apache Spark Deklaratív folyamatok.
Szintaxist és egyéb hivatkozásokat az AUTO CDC INTO (folyamatok), create_auto_cdc_flow és create_auto_cdc_from_snapshot_flow című témakörben talál.
Megjegyzés:
Ez a lap bemutatja, hogyan frissítheti a folyamat tábláit a forrásadatok változásai alapján. A Delta-táblák sorszintű változásadatainak rögzítéséről és lekérdezéséről a Változásadatcsatorna használata Azure Databricks című témakörben olvashat.
Requirements
A CDC API-k használatához a folyamatot úgy kell konfigurálni, hogy kiszolgáló nélküli Lakeflow-folyamatokat vagy Lakeflow-folyamatokat Pro vagy Advancedkiadásokat használjon.
Az AUTO CDC működése
CDC-feldolgozás végrehajtásához AUTO CDC, hozzon létre egy streamelési táblát, majd az AUTO CDC ... INTO utasítást SQL-ben, vagy a create_auto_cdc_flow() függvényt Python-ban használja a forrás, kulcsok és a változási adatfolyam szekvenciájának megadásához. Az adatrögzítés és a pillanatképek módosítása című témakörből megtudhatja, hogyan működik a szekvenálás és az SCD-logika. Tekintse meg az AUTO CDC-példákat.
A változáscsatornával rendelkező forrásból származó kezdeti hidratáláshoz használja AUTO CDC egy once áramlással, majd folytassa a változáscsatorna feldolgozásának folytatását. Lásd : Külső RDBMS-tábla replikálása AUTO CDC használatával.
A szintaxis részleteiért lásd: AUTO CDC INTO (pipelines) vagy create_auto_cdc_flow.
Az AUTO CDC FROM SNAPSHOT működése
AUTO CDC FROM SNAPSHOT a forrásadatok változásait a sorrend szerinti pillanatképek összehasonlításával határozza meg. Pillanatképeket közvetlenül deltatáblából, felhőbeli tárolófájlokból vagy JDBC-ből olvashat. A Python- és az SQL-pipeline felületen is támogatott.
A CDC feldolgozáshoz AUTO CDC FROM SNAPSHOT hozz létre egy streaming táblát, majd definiáld az áramlást:
- Az SQL-ben az
CREATE FLOWutasításAUTO CDC ... FROM SNAPSHOTformáját vagy aCREATE STREAMING TABLEbeágyazottFLOW AUTO CDCzáradékát használja. A forrás egy kötelezőFROM SNAPSHOT (snapshot_query)klauzula, valamint egy opcionálisWITH VERSION (version_query)klauzula, amely kiválasztja a feldolgozandó következő snapshot verziót. Lásd: CREATE FLOW (pipelines). - Python-ban használd a
create_auto_cdc_from_snapshot_flow()függvényt a snapshot, kulcsok és egyéb argumentumok megadására. Lásd create_auto_cdc_from_snapshot_flow.
A két feldolgozási mintázattal és azok használatának időpontjával kapcsolatos részletekért tekintse meg a Snapshot feldolgozási mintázatokat. Tekintse meg az AUTO CDC FROM SNAPSHOT példákat.
Több oszlop használata szekvenáláshoz
Ha több oszlopot (például időbélyeget és azonosítót) szeretne sorrendbe rendezni, használjon egy STRUCT a kombináláshoz. Az API az első mező alapján rendel, döntetlen esetén pedig a második mezőt veszi figyelembe, és így tovább.
SQL
SEQUENCE BY STRUCT(timestamp_col, id_col)
Python
sequence_by = struct("timestamp_col", "id_col")
AUTO CDC példák
Az alábbi példák az SCD 1. és 2. típusú feldolgozását mutatják be változásadatcsatorna-forrás használatával. A mintaadatok új felhasználói rekordokat hoznak létre, törölnek egy felhasználói rekordot, és frissítik a felhasználói rekordokat. Az 1. típusú SCD-példában az utolsó UPDATE műveletek késve érkeznek, és a céltáblából kerülnek ki, ami a rendelésen kívüli eseménykezelést mutatja.
A példákban a következő bemeneti rekordokat használjuk. Ezeket az adatokat úgy hozza létre, hogy a lekérdezést a Mintaadatok létrehozása szakaszban futtatja.
| userId | név | city | művelet | sorszám |
|---|---|---|---|---|
| 124 | Raul | Oaxaca | INSERT | 1 |
| 123 | Isabel | Monterrey | INSERT | 1 |
| 125 | Mercedes | Tijuana | INSERT | 2 |
| 126 | Liliom | Cancun | INSERT | 2 |
| 123 | null | null | töröl | 6 |
| 125 | Mercedes | Guadalajara | UPDATE | 6 |
| 125 | Mercedes | Mexicali | UPDATE | 5 |
| 123 | Isabel | Chihuahua | UPDATE | 5 |
Ha a mintaadat-generálási lekérdezés utolsó sorát törli, a következő rekordot szúrja be, amely megadja a tábla csonkítását (a tábla törlését) a következő helyen sequenceNum=3:
| userId | név | city | művelet | sorszám |
|---|---|---|---|---|
| null | null | null | CSONKÍT | 3 |
Megjegyzés:
Az alábbi példák lehetőségeket tartalmaznak a DELETE és TRUNCATE műveletek megadására, de mindkettő opcionális.
Mintaadatok létrehozása
Mintaadatkészlet létrehozásához futtassa az alábbi utasításokat. Ez a kód nem egy folyamatdefiníció részeként való futtatásra szolgál. Indítsa el a csővezeték kísérleti mappából, az átalakítások mappa helyett.
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.users_cdf
AS SELECT
col1 AS userId,
col2 AS name,
col3 AS city,
col4 AS operation,
col5 AS sequenceNum
FROM (
VALUES
-- Initial load.
(124, "Raul", "Oaxaca", "INSERT", 1),
(123, "Isabel", "Monterrey", "INSERT", 1),
-- New users.
(125, "Mercedes", "Tijuana", "INSERT", 2),
(126, "Lily", "Cancun", "INSERT", 2),
-- Isabel is removed from the system and Mercedes moved to Guadalajara.
(123, null, null, "DELETE", 6),
(125, "Mercedes", "Guadalajara", "UPDATE", 6),
-- This batch of updates arrived out of order. The batch at sequenceNum 6 is the final state.
(125, "Mercedes", "Mexicali", "UPDATE", 5),
(123, "Isabel", "Chihuahua", "UPDATE", 5)
-- Uncomment to test TRUNCATE.
-- ,(null, null, null, "TRUNCATE", 3)
);
1. típusú SCD-frissítések feldolgozása
Az 1. SCD-típus csak az egyes rekordok legújabb verzióját tárolja. Az alábbi példa a fent létrehozott változásadatfolyamból olvas, és módosításokat alkalmaz egy streamelési táblára. Ehhez a kódhoz pipeline kell. Lásd : ETL-folyamatok fejlesztése és hibakeresése a Lakeflow Pipelines-szerkesztővel.
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_current")
dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)
SQL
CREATE OR REFRESH STREAMING TABLE users_current;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_current
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
APPLY AS TRUNCATE WHEN
operation = "TRUNCATE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 1;
Az SCD 1. típusának futtatása után a céltábla a következő rekordokat tartalmazza:
| userId | név | city |
|---|---|---|
| 124 | Raul | Oaxaca |
| 125 | Mercedes | Guadalajara |
| 126 | Liliom | Cancun |
A 123.felhasználó (Izabella) törölve lett, és nem jelenik meg. A 125-ös felhasználó (Mercedes) csak a legújabb várost (Guadalajara) jeleníti meg, mert az 1. scd-típus felülírja a korábbi értékeket. A korábbi UPDATE helyen sequenceNum=5 el lett vetve, mert egy későbbi frissítés sequenceNum=6 megérkezett.
Miután a példát a TRUNCATE rekorddal futtatta, amely nem volt kommentálva, a táblát a rendszer törli a sequenceNum=3 helyen. Ez azt jelenti, hogy a rekordok 124126 nem szerepelnek a táblában, és a végső céltábla csak a következő rekordot tartalmazza:
| userId | név | city |
|---|---|---|
| 125 | Mercedes | Guadalajara |
SCD 2. típusú frissítések feldolgozása
A 2. SCD-típus megőrzi a módosítások teljes előzményeit azáltal, hogy új sorokat hoz létre egy rekord minden verziójához, és __START_AT oszlopokkal __END_AT jelzi, hogy az egyes verziók mikor aktívak.
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2;
A 2. típusú SCD-példa futtatása után a céltábla a következő rekordokat tartalmazza:
| userId | név | city | __START_AT | __END_AT |
|---|---|---|---|---|
| 123 | Isabel | Monterrey | 1 | 5 |
| 123 | Isabel | Chihuahua | 5 | 6 |
| 124 | Raul | Oaxaca | 1 | null |
| 125 | Mercedes | Tijuana | 2 | 5 |
| 125 | Mercedes | Mexicali | 5 | 6 |
| 125 | Mercedes | Guadalajara | 6 | null |
| 126 | Liliom | Cancun | 2 | null |
A táblázat megőrzi a teljes előzményeket. A 123 felhasználó két verzióval rendelkezik (törléskor a 6. sorozatban végződik). A 125-ös felhasználó három verzióval rendelkezik, amelyek városváltozásokat mutatnak. A __END_AT = null rekordok jelenleg aktívak.
SCD 2. típusú oszlop részhalmazának nyomon követése
Alapértelmezés szerint az SCD 2. típusa új verziót hoz létre, amikor bármilyen oszlopérték megváltozik. Megadhatja a nyomon követni kívánt oszlopok egy részhalmazát, így a többi oszlop módosításai nem új előzményrekordot, hanem az aktuális verziót frissítik.
Az alábbi példa kizárja az oszlopot az city előzmények nyomon követéséből:
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2",
track_history_except_column_list = ["city"]
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
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)
Mivel city a módosítások nincsenek nyomon követve, a városfrissítések a jelenlegi sort írják felül új verzió létrehozása helyett. A céltábla a következő rekordokat tartalmazza:
| userId | név | city | __START_AT | __END_AT |
|---|---|---|---|---|
| 123 | Isabel | Chihuahua | 1 | 6 |
| 124 | Raul | Oaxaca | 1 | null |
| 125 | Mercedes | Guadalajara | 2 | null |
| 126 | Liliom | Cancun | 2 | null |
AUTO CDC FROM SNAPSHOT példák
Az alábbi szakaszok példákat mutatnak be a AUTO CDC FROM SNAPSHOT pillanatképek SCD 1. vagy 2. típusú céltáblákba történő feldolgozására. Az API használatának hátteréről az adatrögzítés és a pillanatképek módosítása című témakörben olvashat.
Ebben a szakaszban Python- és SQL-megvalósításokra talál példákat. A teljes SQL-szintaxisért lásd: CREATE FLOW (pipelines).
Példa: Periodikus pillanatképek feldolgozása
Használd ezt a megközelítést, amikor a snapshotok rendszeresen és sorrendben érkeznek. A Python interfész a pipeline frissítési sorrendjét használja a verziókezeléshez. Az SQL esetében a példa külön verziótáblázatot használ annak meghatározására, mikor érhető el új snapshot.
A pillanatképek több forrástípusból is olvashatók, beleértve a Delta-táblákat, a felhőbeli tárolófájlokat és a JDBC-kapcsolatokat.
1. lépés: Mintaadatok létrehozása
Pillanatképadatokat tartalmazó táblázat létrehozása. Futtassa a következő kódot egy notebookból vagy a Databricks SQL-ből a explorations pipeline mappájában:
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.snapshot (
userId INT,
city STRING
);
CREATE TABLE main.cdc_tutorial.snapshot_version (
version BIGINT
);
INSERT INTO main.cdc_tutorial.snapshot VALUES
(1, 'Oaxaca'),
(2, 'Monterrey'),
(3, 'Tijuana');
INSERT INTO main.cdc_tutorial.snapshot_version VALUES (0);
2. lépés: AUTOMATIKUS CDC futtatása PILLANATKÉPBŐL
A lépésben szereplő kód futtatásához pipeline-ra van szükség. Lásd : ETL-folyamatok fejlesztése és hibakeresése a Lakeflow Pipelines-szerkesztővel.
Python esetén válasszon forrástípust a snapshot nézethez (a mintakészítő kód Delta tábla-t generál):
Opció A: Olvasás egy Delta-táblából
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot")
B. lehetőség: Olvasás felhőbeli tárolóból
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.format("csv").option("header", True).load("<snapshot-path>")
C lehetőség: Olvasás JDBC-ből (csak klasszikus számítással)
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)
Ezután add hozzá a céltáblát és a folyamatot. Az SQL implementáció a snapshot_version táblázatot használja annak meghatározására, mikor érhető el új pillanatkép:
Python
dp.create_streaming_table("target")
dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2
)
SQL
CREATE OR REFRESH STREAMING TABLE target;
CREATE FLOW target_snapshot_flow AS
AUTO CDC INTO target
FROM SNAPSHOT (
SELECT userId, city
FROM main.cdc_tutorial.snapshot
)
WITH VERSION (
SELECT version
FROM main.cdc_tutorial.snapshot_version
WHERE (
NOT EXISTS (SELECT 1 FROM last_snapshot_version())
OR version > (SELECT version FROM last_snapshot_version())
)
)
KEYS (userId)
STORED AS SCD TYPE 2;
Az első folyamat futtatása után a rendszer minden rekordot aktív sorként szúr be:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | Oaxaca | 0 | null |
| 2 | Monterrey | 0 | null |
| 3 | Tijuana | 0 | null |
Megjegyzés:
Ha helyette az SCD 1-es típust szeretné használni, és csak az aktuális állapotot szeretné megőrizni, állítsa be a STORED AS SCD TYPE 1 értéket Pythonban vagy a stored_as_scd_type=1 értéket SQL-ben. Ebben az esetben a cél táblázat nem tartalmaz __START_AT és __END_AT oszlopokat.
3. lépés: Új pillanatkép szimulálása és újrafuttatása
Frissítsd a forrástáblát, hogy szimuláld egy új pillanatkép érkezését (futtatd ezt a kódot egy notebookból vagy SQL fájlból a explorations pipeline mappájában):
TRUNCATE TABLE main.cdc_tutorial.snapshot;
INSERT INTO main.cdc_tutorial.snapshot VALUES
(2, 'Carmel'),
(3, 'Los Angeles'),
(4, 'Death Valley'),
(6, 'Kings Canyon');
UPDATE main.cdc_tutorial.snapshot_version SET version = 1;
A folyamat újbóli futtatása.
AUTO CDC FROM SNAPSHOT összehasonlítja az új pillanatképet az előzőhöz, és észleli, hogy az 1. felhasználó törölve lett, a 2. és a 3. felhasználó frissült, és a 4. és a 6. felhasználó lett beszúrva. Ez létrehoz egy változáscsatornát, és a kimeneti tábla létrehozásához használja AUTO CDC .
Az SCD 2. típusával végzett második futtatás után a céltábla a következő rekordokat tartalmazza:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | Oaxaca | 0 | 1 |
| 2 | Monterrey | 0 | 1 |
| 2 | Carmel | 1 | null |
| 3 | Tijuana | 0 | 1 |
| 3 | Los Angeles | 1 | null |
| 4 | Halálvölgy | 1 | null |
| 6 | Kings-kanyon | 1 | null |
Az 1. felhasználó befejeződött (törölve). A 2. és a 3. felhasználó két verzióval rendelkezik, amelyek a város változásait jelenítik meg. A 4. és a 6. felhasználó újonnan lett beszúrva.
Az SCD 1-es típusával végzett második futtatás után a céltábla csak az aktuális állapotot jeleníti meg:
| userId | city |
|---|---|
| 2 | Carmel |
| 3 | Los Angeles |
| 4 | Halálvölgy |
| 6 | Kings-kanyon |
Példa: Pillanatképek feldolgozása verziófüggvényekkel
Ezt a módszert akkor használja, ha explicit módon szabályoznia kell a pillanatképek sorrendjét. Ezt a módszert használhatja például akkor, ha egyszerre több pillanatkép érkezik, vagy a pillanatképek sorrendből érkeznek. Te határozod meg, hogyan lehet kiválasztani a következő snapshotot és annak verziószámát. Az API növekvő verziórendben dolgozza fel a pillanatképeket:
- Ha több pillanatkép is van a tárolóban, az összes feldolgozásuk sorrendben történik.
- Ha egy pillanatkép nem a megfelelő sorrendben érkezik (például
snapshot_3asnapshot_4után érkezik), akkor kihagyják. - Ha nincs új pillanatkép, a verziófüggvény vagy lekérdezés nem ad eredményt, és nem történik feldolgozás.
1. lépés: Pillanatképfájlok előkészítése
Pillanatképadatokat tartalmazó CSV-fájlokat hozhat létre, és hozzáadhatja őket egy kötethez vagy felhőbeli tárolóhelyhez. Nevezze el a fájlokat időrendben (például snapshot_1.csv, snapshot_2.csv).
Minden fájlnak tartalmaznia kell a userId és city oszlopokat. Például:
snapshot_1.csv:
| userId | city |
|---|---|
| 1 | Oaxaca |
| 2 | Monterrey |
| 3 | Tijuana |
snapshot_2.csv:
| userId | city |
|---|---|
| 2 | Carmel |
| 3 | Los Angeles |
| 4 | Halálvölgy |
2. lépés: Az AUTO CDC FROM SNAPSHOT futtatása verziófüggvénnyel
Egy pipeline-ban hozz létre egy új Python vagy SQL fájlt a transformations mappában, és add hozzá a következő kódot. Ezután futtassa a csővezetéket. Lásd : ETL-folyamatok fejlesztése és hibakeresése a Lakeflow Pipelines-szerkesztővel.
Python
from pyspark import pipelines as dp
from typing import Optional, Tuple
from pyspark.sql import DataFrame
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
snapshot_dir = "/Volumes/main/cdc_tutorial/snapshots/" # or the location you created the sample data
files = dbutils.fs.ls(snapshot_dir)
snapshot_files = [f.name for f in files if f.name.startswith("snapshot_") and f.name.endswith(".csv")]
snapshot_versions = []
for filename in snapshot_files:
try:
version = int(filename.replace("snapshot_", "").replace(".csv", ""))
snapshot_versions.append(version)
except ValueError:
continue
snapshot_versions.sort()
if latest_snapshot_version is None:
if snapshot_versions:
next_version = snapshot_versions[0]
else:
return None
else:
next_versions = [v for v in snapshot_versions if v > latest_snapshot_version]
if next_versions:
next_version = next_versions[0]
else:
return None
snapshot_path = f"{snapshot_dir}snapshot_{next_version}.csv"
df = spark.read.format("csv").option("header", True).load(snapshot_path)
return (df, next_version)
dp.create_streaming_table("main.cdc_tutorial.target_versioned")
dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version,
keys = ["userId"],
stored_as_scd_type = 2
)
SQL
CREATE OR REFRESH STREAMING TABLE main.cdc_tutorial.target_versioned;
CREATE FLOW target_versioned_snapshot_flow AS
AUTO CDC INTO main.cdc_tutorial.target_versioned
FROM SNAPSHOT (
SELECT CAST(userId AS INT) AS userId, city
FROM read_files(
'/Volumes/main/cdc_tutorial/snapshots/',
format => 'csv',
header => true
)
WHERE _metadata.file_name = CONCAT(
'snapshot_',
CAST((SELECT version FROM current_snapshot_version()) AS STRING),
'.csv'
)
)
WITH VERSION (
WITH snapshot_files AS (
SELECT
CAST(REGEXP_EXTRACT(path, 'snapshot_([0-9]+)[.]csv$', 1) AS BIGINT) AS version
FROM list_files('/Volumes/main/cdc_tutorial/snapshots/')
WHERE path RLIKE 'snapshot_[0-9]+[.]csv$'
)
SELECT version
FROM snapshot_files
WHERE (
NOT EXISTS (SELECT 1 FROM last_snapshot_version())
OR version > (SELECT version FROM last_snapshot_version())
)
ORDER BY version
LIMIT 1
)
KEYS (userId)
STORED AS SCD TYPE 2;
Megjegyzés:
Ha ehelyett az SCD 1-es típust szeretné használni, állítsa be a STORED AS SCD TYPE 1 értéket Pythonban vagy a stored_as_scd_type=1 értéket SQL-ben.
A feldolgozás snapshot_1.csvután a céltábla a következő rekordokat tartalmazza:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | Oaxaca | 1 | null |
| 2 | Monterrey | 1 | null |
| 3 | Tijuana | 1 | null |
A feldolgozás snapshot_2.csvután a céltábla a következő rekordokat tartalmazza:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | Oaxaca | 1 | 2 |
| 2 | Monterrey | 1 | 2 |
| 2 | Carmel | 2 | null |
| 3 | Tijuana | 1 | 2 |
| 3 | Los Angeles | 2 | null |
| 4 | Halálvölgy | 2 | null |
Megjegyzés:
Ne feledje, hogy az 1. SCD-típus esetében a táblázat pontosan úgy néz ki, mint a legújabb pillanatkép. A különbség az, hogy az alsóbb rétegbeli lekérdezések csak a módosított rekordok feldolgozására használhatják a változáscsatornát.
3. lépés: Új pillanatképek hozzáadása
Adjon hozzá egy új CSV-fájlt a tárolóhelyhez módosított adatokkal (például módosított városértékekkel, új sorokkal vagy eltávolított sorokkal). Ezután futtassa újra a folyamatot az új pillanatkép feldolgozásához.
Korlátozások
- A szekvenáló oszlopnak rendezhető adattípusnak kell lennie.
NULLa szekvenálási értékek nem támogatottak. - Ha adatokat szeretne streamelni egy AUTO CDC-folyamat célpontjából, olvassa ki azokat a változási hírcsatornából. Részletekért lásd: Változásadatcsatorna olvasása AUTO CDC céltáblából.
További erőforrások
- Adatrögzítés és pillanatképek módosítása: Tudnivalók a CDC-fogalmakról, a pillanatképekről és az SCD-típusokról.
-
Egy külső RDBMS-tábla replikálása a következővel
AUTO CDC: Ismerje meg, hogyan végezhet kezdeti hidratálást egyoncefolyamattal, majd folytathatja a módosítások feldolgozását. - Speciális AUTO CDC-témakörök: Ismerje meg az AUTOMATIKUS CDC-célok változási műveleteit, a változásadatcsatornák olvasását és a metrikák feldolgozását.
- Tekerd vissza és játszd le egy csővezetéket: Tanuld meg, hogyan állíthatod vissza egy AUTO CDC célpontot egy korábbi időpontba, és újraértelmezheted a változtatásokat.
- Oktatóanyag: ETL-folyamat létrehozása változásadat-rögzítéssel