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.
RDBMS-tábla replikálása külső forrással
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
oncefolyamattal. - Új módosítások folyamatos beolvasása egy
changefolyamattal.
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 CDCfolyamat 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.readStreamautomatikus betöltővel való használata (format("cloudFiles")) - JSON-fájlokat olvas be egy olyan könyvtárból, amelyet
orders_snapshot_path - Az
includeExistingFilesbeállításatrueá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, hogytrueautomatikusan 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
- Teljes körűen felügyelt SQL Server-összekötő a Lakeflow Connectben