RDBMS-tábla replikálása külső forrással AUTO CDC

Egy táblát egy külső relációsadatbázis-kezelő rendszerből (RDBMS) replikálhat az Azure Databricksbe a folyamatokban a AUTO CDC API használatával. A következőt fogja elsajátítani:

  • A források beállításának gyakori mintái.
  • Hogyan hajtsunk végre egy egyszeri teljes másolatot a meglévő adatokról egy once folyamattal.
  • Új módosítások folyamatos beolvasása egy change folyamattal.

Ez a minta ideális a lassan változó dimenziótáblák (SCD) létrehozásához, vagy egy céltábla külső rekordrendszerrel való szinkronizálásához.

Mielőtt hozzákezdene

Ez az útmutató feltételezi, hogy a forrásból a következő adathalmazokhoz fér hozzá:

  • A forrástábla teljes pillanatképe a felhőbeli tárolóban. Ez az adatkészlet a kezdeti betöltéshez használatos.
  • Folyamatos változáscsatorna, amely ugyanazon a felhőbeli tárolóhelyen van feltöltve (például Debezium, Kafka vagy naplóalapú CDC használatával). Ez a hírcsatorna a folyamatban lévő AUTO CDC folyamat bemenete.

Forrásnézetek beállítása

Először határozzon meg két forrásnézetet, amelyek a rdbms_orders céltáblát a felhőtárolási útvonal orders_snapshot_path alapján töltik fel. Mindkettő a felhőtárolóban található nyers adatokra épülő stream nézetként van kialakítva. A nézetek használata nagyobb hatékonyságot biztosít, mivel az adatokat nem kell megírni a AUTO CDC folyamat során való használat előtt.

  • Az első forrásnézet egy teljes pillanatkép (full_orders_snapshot)
  • A második egy folyamatos változási hírcsatorna (rdbms_orders_change_feed).

Az útmutatóban szereplő példák a felhőbeli tárolást használják forrásként, de a streamelési táblák által támogatott forrásokat is használhatja.

full_orders_snapshot()

Ez a lépés létrehoz egy adatcsővezetéket egy nézettel, amely beolvassa a rendelések adatainak kezdeti teljes állapotfelvételét.

Python

A következő Python-példa:

  • Az spark.readStream automatikus betöltővel való használata (format("cloudFiles"))
  • JSON-fájlokat olvas be egy olyan könyvtárból, amelyet orders_snapshot_path
  • Az includeExistingFiles beállítása true állapotra annak érdekében történik, hogy az útvonalon már meglévő előzményadatok feldolgozásra kerüljenek.
  • Beállítja inferColumnTypes, hogy true automatikusan következtetni tudjon a sémára.
  • Minden oszlopot visszaadja .select("\*") szerint
@dp.view()
def full_orders_snapshot():
    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.includeExistingFiles", "true")
        .option("cloudFiles.inferColumnTypes", "true")
        .load(orders_snapshot_path)
        .select("*")
    )

SQL

A következő SQL-példa ugyanazokat a beállításokat név szerinti argumentumokként adja át a read_files számára. orders_snapshot_path SQL-változóként kell rendelkezésre állnia (például folyamatparaméterekkel definiálva vagy manuálisan interpolálva).

CREATE OR REPLACE TEMPORARY VIEW full_orders_snapshot
AS SELECT *
FROM STREAM read_files("${orders_snapshot_path}",
  format => "json",
  includeExistingFiles => "true",
  inferColumnTypes => "true"
);

rdbms_orders_change_feed()

Ez a lépés létrehoz egy második nézetet, amely beolvassa a növekményes változásadatokat (például CDC-naplókból vagy változástáblákból). A orders_cdc_path olvassa és feltételezi, hogy a CDC-stílusú JSON-fájlok rendszeresen ebbe az elérési útra kerülnek.

Python

@dp.view()
def rdbms_orders_change_feed():
  return (
    spark.readStream
    .format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.includeExistingFiles", "true")
    .option("cloudFiles.inferColumnTypes", "true")
    .load(orders_cdc_path)
  )

SQL

Az alábbi SQL-példában egy változó látható, ${orders_cdc_path} amely interpolálható a folyamatbeállítások egyik értékének beállításával, vagy egy változó explicit beállításával a kódban.

CREATE OR REPLACE TEMPORARY VIEW rdbms_orders_change_feed
AS SELECT *
FROM STREAM read_files("${orders_cdc_path}",
  format => "json",
  includeExistingFiles => "true",
  inferColumnTypes => "true"
);

Kezdeti hidratálás (egyszeri áramlás)

A források beállítása AUTO CDC után a logika mindkét forrást egy célstreamelési táblába egyesíti. Először használjon egyszeri AUTO CDC folyamatot ONCE=TRUE az RDBMS-tábla teljes tartalmának streamelő táblába másolásához. Ezzel előkészíti a céltáblát az előzményadatokkal anélkül, hogy a későbbi frissítések során újra végrehajtaná azokat.

Python

from pyspark import pipelines as dp

# Step 1: Create the target streaming table

dp.create_streaming_table("rdbms_orders")

# Step 2: Once Flow — Load initial snapshot of full RDBMS table

dp.create_auto_cdc_flow(
  flow_name = "initial_load_orders",
  once = True,  # one-time load
  target = "rdbms_orders",
  source = "full_orders_snapshot",  # e.g., ingested from JDBC into bronze
  keys = ["order_id"],
  sequence_by = "timestamp",
  stored_as_scd_type = "1"
)

SQL


-- Step 1: Create the target streaming table
CREATE OR REFRESH STREAMING TABLE rdbms_orders;

-- Step 2: Once Flow for initial snapshot
CREATE FLOW rdbms_orders_hydrate
AS AUTO CDC ONCE INTO rdbms_orders
FROM stream(full_orders_snapshot)
KEYS (order_id)
SEQUENCE BY timestamp
STORED AS SCD TYPE 1;

A once folyamat csak egyszer fut. A csoportfolyamat full_orders_snapshot létrehozása után hozzáadott új fájlokat nem vesszük figyelembe.

Fontos

A streamelési tábla teljes frissítésének rdbms_orders végrehajtása újrafuttatja a once folyamatot. Ha a kezdeti pillanatképadatok a felhőbeli tárolóban el lettek távolítva, ez adatvesztést eredményez.

Folyamatos változási hírcsatorna (változásfolyam)

A kezdeti pillanatkép betöltése után használjon egy másik AUTO CDC folyamatot az RDBMS CDC-hírcsatornájának változásainak folyamatos betöltéséhez. Ez naprakészen tartja a rdbms_orders táblázatot beszúrásokkal, frissítésekkel és törlésekkel.

Python

from pyspark import pipelines as dp

# Step 3: Change Flow — Ingest ongoing CDC stream from source system

dp.create_auto_cdc_flow(
flow_name = "orders_incremental_cdc",
target = "rdbms_orders",
source = "rdbms_orders_change_feed", # e.g., ingested from Kafka or Debezium
keys = ["order_id"],
sequence_by = "timestamp",
stored_as_scd_type = "1"
)

SQL

-- Step 3: Continuous CDC ingestion
CREATE FLOW rdbms_orders_continuous
AS AUTO CDC INTO rdbms_orders
FROM stream(rdbms_orders_change_feed)
KEYS (order_id)
SEQUENCE BY timestamp
STORED AS SCD TYPE 1;

Megfontolások

Idempotencia visszatöltése A once folyamat csak a céltábla teljes frissítésekor fut újra.
Több folyam Több változási folyamat használatával egyesítheti a javításokat, a későn érkező adatokat vagy az alternatív hírcsatornákat, de mindegyiknek meg kell osztania egy sémát és kulcsokat.
Teljes frissítés A teljes frissítés az rdbms_orders adatfolyam táblán újraindítja a once adatfolyamot. Ez adatvesztéshez vezethet, ha a kezdeti felhőbeli tárolóhely törölte a kezdeti pillanatképadatokat.
Folyamat végrehajtási sorrendje A folyamat végrehajtásának sorrendje nem számít. A végeredmény ugyanaz.

További erőforrások