Az AUTO CDC API-k: Az adatrögzítés egyszerűsítése csővezetékekkel

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 CDC egy 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 FLOW utasítás AUTO CDC ... FROM SNAPSHOT formáját vagy a CREATE STREAMING TABLE beágyazott FLOW AUTO CDC záradékát használja. A forrás egy kötelező FROM SNAPSHOT (snapshot_query) klauzula, valamint egy opcionális WITH 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_3 a snapshot_4 utá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. NULL a 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