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.
Important
Ez a funkció bétaverzióban érhető el.
A Databricks Python egységteszteltségével kapcsolatos általános információkért lásd Python egységtesztelést.
A Lakeflow-folyamatok támogatják Python egységtesztek írását a webes Lakeflow Pipelines-szerkesztőben. Ez lehetővé teszi Python vagy SQL-átalakítási logika érvényesítését a modelladatok használatával. A folyamattesztelési keretrendszerrel tesztelheti a peremhálózati eseteket, érvényesítheti a saját fejlesztésű folyamat API-kat (automatikus CDC, streamelési táblák, elvárások, hozzáfűző folyamatok), és iterálhat a támogatott táblaazonosító műveletekhez használt mintabemenetek használatával. A tesztek futtatása előtt tekintse át az elkülönítési korlátozásokat.
- Izolált tesztvégrehajtás: A keretrendszer egy SparkSession-t biztosít, amely átirányítja a táblaműveleteket a folyamat alapértelmezett katalógusában lévő ideiglenes tesztsémára, így a bemeneti adatokat ki lehet modellezni és tesztkimeneteket írni anélkül, hogy az hatással van az éles táblákra. Az elkülönítés olyan műveletekre vonatkozik, amelyek név szerint hivatkoznak egy táblára; lásd : Korlátozások.
- Rugalmas tesztelési hatókör: A teszt SparkSession használatával hajtsa végre egy folyamatcső részhalmazát (egyes táblákat, egymástól függő táblák láncait vagy teljes folyamatcsöveket) a folyamatcső számítási erőforrásain.
- Eredményérvényesítés: A tesztben létrehozott izolált kimeneti táblák eredményeinek ellenőrzése standard pytest-állításokkal.
Mikor érdemes egységtesztelést használni?
A tipikus használati esetek a következők:
- Új átalakítási logika érvényesítése: Tesztelje, hogy az átalakítás létrehozza-e a várt sémát, a sorok számát, az összesítéseket és az üzleti logikát, mielőtt éles adatokon fut.
- Automatikus CDC-specifikációk tesztelése: Ellenőrizze, hogy az automatikus CDC-folyamatdefiníciók helyesen dolgozzák-e fel a módosítási eseményeket, kezelik-e a beszúrásokat, frissítéseket, törléseket és SCD-típusokat (lassan módosítva a dimenziót) a modelladatok használatával.
- Elvárások és adatminőségi szabályok tesztelése: Ellenőrizze, hogy az elvárások akkor vallanak kudarcot, amikor kell, és sikeresen teljesülnek, amikor az adatok érvényesek.
- Tesztelés függő táblákon: Tesztelje az átalakítások láncait (például bronz, ezüst és arany), hogy ellenőrizze, hogy az adatok megfelelően haladnak-e végig a folyamatdiagramon.
Requirements
A folyamatengedélyek
Owner, valamint aUSE CATALOGCREATE SCHEMAfolyamat alapértelmezett katalógusában szereplő jogosultságok. A keretrendszernek szüksége van ezekre a jogosultságokra az ideiglenes tesztséma létrehozásához, ahol a tesztek futnak.A folyamat engedélyének ellenőrzéséhez vagy beállításához nyissa meg a folyamatot, és kattintson a Megosztás gombra. Önnek kell lennie a pipeline-nak
Owner(IS OWNER); aCAN RUNésCAN MANAGEnem elegendők a tesztek futtatásához. Lásd: Folyamatengedélyek konfigurálása.A katalógus jogosultságainak ellenőrzéséhez vagy beállításához nyissa meg a katalógust a Katalóguskezelőben, válassza az Engedélyek lapot, és győződjön meg arról, hogy rendelkezik
USE CATALOGésCREATE SCHEMA. A katalógus tulajdonosa, a metaadattár rendszergazdája vagy a jogosultsággal rendelkező felhasználó a következő jogosultságokkal ruházhatja fel őket, beleértve azMANAGESQL-t is:GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;További információért tekintse meg a Unity Catalog jogosultságok referenciája részt.
A folyamatcsatornát eseményindítású (nem folyamatos) módra kell konfigurálni.
A pipeline-nak a PREVIEW csatornán kell lennie. Az egységtesztelés bétaverzióban érhető el, és csak előzetes verzióban érhető el.
A Spark Connect nem támogatott.
Note
A tesztelkülönítés olyan táblaműveleteket fed le, amelyek név szerint hivatkoznak egy táblára. Az elkülönítést megkerülő műveletek a tesztkódban és a kiválasztott kimenetek által végrehajtott folyamatkódokban is előfordulhatnak, beleértve azok tranzitív függőségeit is. Egy biztonságosnak tűnő tesztfájl továbbra is futtathat olyan pipeline-folyamatot, amely elérési út vagy csatlakozó használatával olvas vagy ír, és amely az éles adatokon végez műveleteket. Annak érdekében, hogy a tesztek ne legyenek hatással az éles adatokra vagy metaadatokra, kövesse az alábbi szabályokat:
- Hivatkozzon minden táblára név szerint (
catalog.schema.table), és mockoljon minden bemenetet név szerint. Ne olvasson és ne írjon elérési út alapján (/Volumes/...,dbfs:/...,s3://...,abfss://...), és ne olvasson olyan konnektorokból, mint a Kafka vagy az Auto Loader. Ezek megkerülik az elkülönítést, és közvetlenül valódi éles rendszerekre hatnak. - Ne futtasson irányítási vagy tulajdonosi utasításokat, például
GRANT: ,REVOKEALTER ... OWNER TO,SET/UNSET TAGSvagy .CREATE/DROP POLICYEzek a valódi, biztonságos éles környezetben hajthatók végre. - Ne hozzon létre katalógusokat vagy sémákat (
CREATE CATALOG,CREATE SCHEMA). Ezek elérik a valódi Unity Catalog-metaadattárat. - Ne futtassa a teljes folyamatot, ha a gráf elérésiút-alapú bemeneteket, összekötőket, imperatív írásokat vagy egyéb külső mellékhatásokat tartalmaz. Csak azokat a kimeneteket válassza ki, amelyek függőségei támogatott katalógustáblázat-műveleteket használnak, és amelyeket a rendszer modellbemenetekre cserélt.
További részletekért lásd : Korlátozások .
Limitations
Warning
Egyes műveletek megkerülik a tesztelkülönítést, és valós üzemi adatokra vagy metaadatokra is képesek. A tesztek futtatása előtt tekintse át az alábbi korlátozásokat.
A tesztek elkülönítése kizárólag a tábla neve alapján történik
Ne olvasson és ne írjon elérési út vagy összekötő alapján. Az elkülönítés csak azokat a műveleteket irányítja át, amelyek név alapján hivatkoznak egy táblára (például
spark.read.table("catalog.schema.table")).df.write.saveAsTable("catalog.schema.table")Az útvonalon vagy összekötőn keresztül elért műveletek megkerülik az izolációt, és közvetlenül a tényleges éles rendszereken hajtódnak végre:-
Az elérési útra történő írás (például
df.write.save("/Volumes/..."), egydbfs:/elérési útra, vagy egy felhőbeli vagy külső helyre mutató elérési útra, példáuls3://...vagyabfss://...) a tényleges éles tárolóra ír, és felülírhatja az éles adatokat. -
Az elérési út szerinti olvasás (például
spark.read.load(path)spark.read.format("delta").load(path)) valós üzemi adatokat ad vissza a modell helyett. - A csatlakozóból történő olvasás a tényleges éles forráshoz csatlakozik. Ebbe beletartozik a Kafka (a valódi közvetítőktől származó olvasások) és az Automatikus betöltő (
cloudFilesamely a valódi felhőbeli tároló elérési útjáról olvas). Egyik sem kerül átirányításra a tesztadatokra.
-
Az elérési útra történő írás (például
Ne használja a
event_log()táblaértékelt függvényt folyamategység-tesztből. Teszt módban a(z)event_log()nincs átirányítva a tesztfuttatásának eseménynaplójába. Visszaadhatja a produkciós vagy egy korábban rögzített eseménynaplót, így az erre vonatkozó aszerciók produkciós adatokat olvashatnak. Ehelyett használja a futtatás által visszaadottevent_log_table_nameelemet, és a(z)test_sparkhasználatával kérdezze le.event_log_table_namelehetNone(például ha az eseménynapló tábla neve nem oldható fel), ezért ellenőrizze a lekérdezés előtt:status = test_pipeline.run(test_spark, set(["catalog.schema.table"])) assert status.event_log_table_name is not None events = test_spark.table(status.event_log_table_name)Ne érvényesítse
status.is_successaz eseménynapló elolvasása előtt, ha a cél egy sikertelen frissítés diagnosztizálása. Az eseménynapló gyakran az, amit megvizsgál, hogy megértse, miért hiúsult meg egy frissítés.
Irányítási és DDL-műveletek
- A katalógus, séma, engedély, tulajdonjog, címke és szabályzatmutációk nem támogatottak. Ebbe beletartozik
CREATE/DROP/ALTER CATALOGa (CREATE/DROP/ALTER SCHEMAbeleértve ),SET MANAGED LOCATIONGRANT/REVOKE,ALTER ... OWNER TOSET/UNSET TAGSés .CREATE/DROP POLICYAtest_sparkelemen keresztül végrehajtott SQL-utasítások egyes formáit a rendszer a többrétegű védelem részeként elutasítja; más formák, illetve ugyanezek a műveletek közvetlen API-kon keresztül meghívva elérhetik a tényleges éles objektumokat. Ne hagyatkozz ezekre az őrökre elkülönítési határként. Ezeket az utasításokat tartsa ki a tesztkódból és a kiválasztott kimenetek által végrehajtott folyamatkódból.
Működési korlátozások
- Az egyidejű végrehajtás nem támogatott: A teszt és a folyamatfrissítés egyidejű futtatása nem támogatott, és a rendszer nem akadályozza meg. A kettő között nincs koordináció, így az egyidejű futtatás az erőforrásokért is küzdhet, jelentősen ronthatja az éles frissítés teljesítményét, vagy a teszt elindítása meghiúsulhat. Ne indítsa el a tesztet, amíg a folyamat frissítést futtat (vagy ne indítson frissítést egy teszt futtatása közben); várjon, amíg a folyamatban lévő frissítések befejeződnek a tesztek futtatása előtt.
-
Ideiglenes sémák rendellenes leállítás után: Minden tesztfuttatás létrehoz egy ideiglenes sémát (névvel elnevezve
redirecting_<id>) a folyamat alapértelmezett katalógusában, és automatikusan elveti, amikor a futtatás befejeződik. Ha egy futtatás rendellenesen végződik (például a számítás futás közben elveszik), az ideiglenes séma hátrahagyható, a futtatás minta- és kimeneti tábláit tartva. Ez nem befolyásolja az éles adatokat. A tárterület visszaállításához manuálisan helyezze el azokat a hátrahagyott sémákat, amelyek neve a folyamat alapértelmezett katalógusában kezdődikredirecting_. - A tesztfuttatások számítást használnak: A tesztfuttatások a folyamat számításán futnak, és a számlázás normál folyamatfrissítésként történik. A tesztfuttatásokhoz nincs külön mérési lehetőség.
-
A teljes frissítés nem támogatott: Csak szelektív frissítés érhető el.
test_pipeline.run()frissíti a kiválasztott kimeneteket (vagy kijelölés hiányában az összes kimenetet); a teljes frissítés és a teljes frissítés kijelölése nincsenek megvalósítva.
Szerzői és hűségkorlátok
- Csak szerkesztői végrehajtás: A teszteket a webes Lakeflow Pipelines-szerkesztőből kell futtatni.
- csak Python tesztek: A teszteket Python kell írni. Tesztelheti az SQL-adatfolyamokat, de magukat a teszteket Pythonban kell megírni.
- Irányítási hűség: A modelladatok nem öröklik a lecserélt éles táblákon definiált sorszűrőket vagy oszlopmaszkokat. A teszteredmények pontosan az Ön által megadott mintabemeneteket tükrözik, és eltérhetnek attól, hogy ugyanaz a lekérdezés hogyan viselkedik a szabályozott éles adatokon.
1. lépés: Folyamatbeállítások frissítése
Konfigurálja a folyamatot úgy, hogy aktivált módban fusson a PREVIEW csatornán.
- A felhasználói felületen nyissa meg a folyamatot, és kattintson a Beállítások>speciális beállításai>csatorna>előnézete elemre
- Állítsa be a folyamat üzemmódjátaktiváltra (ne használja a folyamatos üzemmódot).
Másik lehetőségként szerkessze közvetlenül a folyamat beállításait tartalmazó JSON-fájlt:
"continuous": false,
"channel": "PREVIEW"
2. lépés: Tesztfájl létrehozása
A Lakeflow Pipelines-szerkesztőben kattintson a + (hozzáadás) gombra, és válassza a Tesztelés lehetőséget. Ez létrehoz egy tesztfájlt (és tests a mappát, ha még nem létezik), amely nem szerepel a folyamat forráskódjában. Nem kell saját maga létrehoznia a tests mappát.
3. lépés: Tesztek létrehozása
A Genie Code képes tesztállványok létrehozására:
A tesztfájlban kattintson a Tesztek létrehozása gombra.
Másik lehetőségként használja a
/testselemet a Genie Code ügynöki módjában.
Használja a Genie Code-ot sablonkód létrehozására, majd szabja testre a szélső esetekhez.
Másik lehetőségként saját maga is megírhatja a tesztkódot. Adja hozzá a következő importálásokat az egyes tesztfájlok elejéhez:
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
4. lépés: Tesztek futtatása
Teszteket futtathat a Lakeflow Pipelines-szerkesztőből:
- Kattintson a
(lejátszás) gombra a tesztfüggvény melletti margón egy egyéni teszt futtatásához.
- Kattintson a Tesztfájlok futtatása fájlban elemre a tesztfájl tetején az adott fájlban lévő összes teszt futtatásához.
A tesztelési eredmények (sikeres vagy sikertelen) megjelennek a Szerkesztő alsó paneljén. Tekintse át az assertálási hibákat a sikertelen futások hibakereséséhez.
API-k tesztelése
| API | Description |
|---|---|
TestPipeline.active() |
A Lakeflow Pipelines Editorban éppen szerkesztett folyamathoz egy TestPipeline objektumot ad vissza. Ez az objektum a folyamatra mutató hivatkozás, beleértve annak forráskódját, konfigurációit, alapértelmezett katalógusát/sémáját stb. |
test_pipeline.run(test_spark, set([table_names])) |
Szinkron módon végrehajtja a folyamat frissítését, és a táblanevek megadása esetén szelektív frissítést hajt végre. A folyamatcső végrehajtásának sikeres befejeződése vagy kivétellel történő megszakadása után tér vissza. |
test_spark rögzítőelem |
Létrehoz egy teszt SparkSessiont katalógus–tábla átirányítással, amely automatikusan egy ideiglenes tesztsémára irányítja át azokat a táblaolvasási és -írási műveleteket, amelyek egy táblára név alapján hivatkoznak (például spark.read.table("catalog.schema.table") vagy df.write.saveAsTable("catalog.schema.table")). Az átirányítás csak a névalapú táblaműveletekre vonatkozik; nem terjed ki az elérési úton vagy összekötőn keresztül megcímzett olvasásokra és írásokra, amelyek közvetlenül a valós rendszeren működnek. Lásd Korlátozások. |
Mintaadatok létrehozása
A bemeneti adatokat SQL-lel vagy a(z) createDataFrame használatával szimulálhatja:
# Option 1: Using SQL
test_spark.sql("""
CREATE TABLE catalog.schema.table_name AS
SELECT * FROM VALUES
(1, 'value1'),
(2, 'value2')
AS t(id, name)
""")
# Option 2: Using createDataFrame
df = test_spark.createDataFrame(
[(1, 'value1'), (2, 'value2')],
schema=["id", "name"]
)
df.write.saveAsTable("catalog.schema.table_name")
A valósághű szintetikus adatok nagyobb mennyiségének létrehozásához használhatja a Faker-kódtárat. Először futtassa a(z) %pip install faker elemet a folyamatban, majd hozzon létre egy DataFrame-et Faker-alapú UDF-ekből:
# Option 3: Using Faker for synthetic data
from pyspark.sql import functions as F
from faker import Faker
fake = Faker()
fake_firstname = F.udf(fake.first_name)
fake_lastname = F.udf(fake.last_name)
fake_email = F.udf(fake.ascii_company_email)
df = (
test_spark.range(0, 100)
.withColumn("firstname", fake_firstname())
.withColumn("lastname", fake_lastname())
.withColumn("email", fake_email())
)
df.write.saveAsTable("catalog.schema.table_name")
A folyamatlánc vagy a megadott táblák futtatása
# Run specific tables
test_pipeline.run(test_spark, set(["catalog.schema.table1", "catalog.schema.table2"]))
# Run all tables in the pipeline
test_pipeline.run(test_spark)
Examples
1. példa: A sorok számával, sémájával és null kezelésével rendelkező összesítések tesztelése
Cél: A felhasználói összesítés ellenőrzése helyesen számlálja meg a felhasználókat típus szerint, kezeli a null e-maileket, és létrehozza a várt sémát.
Folyamatátalakítások:
Ezek az átalakítások egy egyszerű kéttáblás folyamatot hoznak létre: users kiválasztja a felhasználói adatokat, és counts típus szerint csoportosítja a felhasználókat, és megszámolja az összes felhasználót és az érvényes e-maileket.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, count, count_if
@dp.table
def users():
return (
spark.read.table("catalog.schema.wanderbricks_users")
.select("user_id", "email", "name", "user_type")
)
@dp.table
def counts():
return (
spark.read.table("catalog.schema.users")
.withColumn("valid_email", col("email").isNotNull())
.groupBy("user_type")
.agg(
count("user_id").alias("total_count"),
count_if("valid_email").alias("count_valid_emails")
)
)
Tesztek:
Ezek a tesztek ellenőrzik a sorszámokat, a sémastruktúrát, a nullkezelést és az összesítési logikát azáltal, hogy szándékosan null értékekkel hoznak létre felhasználói adatokat, és elkülönítve futtatják a folyamatot.
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
from pyspark.testing import assertDataFrameEqual
test_pipeline = TestPipeline.active()
# Mock data fixture
def mock_users(session):
session.sql("""
CREATE TABLE catalog.schema.wanderbricks_users AS
SELECT * FROM VALUES
(1, 'alice@example.com', 'Alice', 'admin'),
(2, NULL, 'Bob', 'user'),
(3, 'charlie@example.com', 'Charlie', 'user'),
(4, NULL, 'Dana', 'admin')
AS t(user_id, email, name, user_type)
""")
# Test 1: Row count
def test_users_row_count(test_spark):
mock_users(test_spark)
test_pipeline.run(test_spark, set(["catalog.schema.users"]))
result = test_spark.table("catalog.schema.users")
assert result.count() == 4
# Test 2: Schema validation
def test_users_schema(test_spark):
mock_users(test_spark)
test_pipeline.run(test_spark, set(["catalog.schema.users"]))
result = test_spark.table("catalog.schema.users")
expected_fields = {"user_id", "email", "name", "user_type"}
actual_fields = set(f.name for f in result.schema.fields)
assert expected_fields == actual_fields
# Test 3: Null handling
def test_users_null_handling(test_spark):
mock_users(test_spark)
test_pipeline.run(test_spark, set(["catalog.schema.users"]))
result = test_spark.table("catalog.schema.users")
null_emails = result.filter("email IS NULL").count()
assert null_emails == 2
# Test 4: Aggregation
def test_counts(test_spark):
mock_users(test_spark)
# Run both tables since counts depends on users
test_pipeline.run(test_spark, set(["catalog.schema.users", "catalog.schema.counts"]))
result = test_spark.table("catalog.schema.counts")
# Check counts for each user_type
admin_row = result.filter("user_type = 'admin'").collect()[0]
user_row = result.filter("user_type = 'user'").collect()[0]
assert admin_row["total_count"] == 2
assert admin_row["count_valid_emails"] == 1
assert user_row["total_count"] == 2
assert user_row["count_valid_emails"] == 1
# Test 5: Full DataFrame comparison with assertDataFrameEqual
def test_counts_full_dataframe(test_spark):
mock_users(test_spark)
test_pipeline.run(test_spark, set(["catalog.schema.users", "catalog.schema.counts"]))
result = test_spark.table("catalog.schema.counts")
expected = test_spark.createDataFrame(
[("admin", 2, 1), ("user", 2, 1)],
schema=["user_type", "total_count", "count_valid_emails"]
)
assertDataFrameEqual(result, expected)
2. példa: Az automatikus CDC tesztelése
Cél: Annak ellenőrzése, hogy az Auto CDC helyesen dolgozza fel a beszúrásokat és frissítéseket tartalmazó változáscsatornát.
Folyamatátalakítás:
Ez az átalakítás beállítja az automatikus CDC-t egy változáscsatornából, amely beolvassa a streamelési módosításokat, és SCD 1-es típusként alkalmazza őket a céltáblára (csak a legújabb verziót tartja meg).
from pyspark import pipelines as dp
from pyspark.sql.functions import col
@dp.view
def users():
return spark.readStream.table("catalog.schema.change_feed")
dp.create_streaming_table("target_autocdc")
dp.create_auto_cdc_flow(
target="target_autocdc",
source="users",
keys=["userId"],
sequence_by=col("ts"),
stored_as_scd_type=1
)
Tesztek:
Az első teszt egy szimulált változáscsatornát hoz létre, amely ugyanahhoz a userId-hoz több rekordot tartalmaz (egy frissítés szimulálására), és ellenőrzi, hogy a céloldalon csak a legfrissebb rekord marad meg. A második teszt a késve érkező és a nem sorrendben érkező eseményeket a folyamat futtatásával, további eseményeknek a változáscsatornához való hozzáfűzésével, majd a folyamat újbóli futtatásával szimulálja.
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
# Test 1: Standard inserts and updates
def test_auto_cdc_flow(test_spark):
# Create a mock change feed table
test_spark.sql("""
CREATE TABLE catalog.schema.change_feed AS
SELECT * FROM VALUES
(1, 'Alice', 1000),
(2, 'Bob', 1001),
(1, 'Alice Updated', 1002)
AS t(userId, name, ts)
""")
# Run the pipeline
test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))
# Read the output
result = test_spark.table("catalog.schema.target_autocdc")
# Verify two users exist
user_ids = set(row["userId"] for row in result.collect())
assert user_ids == {1, 2}
# Verify latest record for userId=1 has ts=1002
latest_user1 = result.filter("userId = 1").collect()[0]
assert latest_user1["ts"] == 1002
assert latest_user1["name"] == "Alice Updated"
# Verify userId=2 has ts=1001
user2 = result.filter("userId = 2").collect()[0]
assert user2["ts"] == 1001
# Test 2: Late-arriving and out-of-order events
def test_auto_cdc_late_arriving(test_spark):
# First batch of change events
test_spark.sql("""
CREATE TABLE catalog.schema.change_feed AS
SELECT * FROM VALUES
(1, 'Alice', 1000),
(2, 'Bob', 1001)
AS t(userId, name, ts)
""")
# Run the pipeline with the initial batch
test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))
# Append late-arriving events to the change feed:
# - A newer event for userId=1 (ts=1003) that arrived after the first run
# - A stale event for userId=2 (ts=999) with a timestamp older than what is already applied
test_spark.sql("""
INSERT INTO catalog.schema.change_feed VALUES
(1, 'Alice Updated', 1003),
(2, 'Bob (stale)', 999)
""")
# Re-run the pipeline. sequence_by=ts ensures stale events do not overwrite newer state.
test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))
result = test_spark.table("catalog.schema.target_autocdc")
# userId=1 should reflect the newer late-arriving event
alice = result.filter("userId = 1").collect()[0]
assert alice["ts"] == 1003
assert alice["name"] == "Alice Updated"
# userId=2 should be unchanged: the stale event with an older ts is ignored
bob = result.filter("userId = 2").collect()[0]
assert bob["ts"] == 1001
assert bob["name"] == "Bob"
3. példa: Automatikus CDC tesztelése pillanatképből
Cél: Ellenőrizze, hogy a CDC megfelelően dolgozza-e fel a pillanatkép módosításait, beleértve a beszúrásokat, frissítéseket és törléseket.
Folyamatátalakítás:
Ez az átalakítás beállítja az automatikus CDC-t a pillanatképből, amely egy pillanatképtáblából olvas, és az idő függvényében követi nyomon a változásokat SCD 2-es típusként (teljes előzményt tart fenn).
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("catalog.schema.snapshot")
dp.create_streaming_table("catalog.schema.target")
dp.create_auto_cdc_from_snapshot_flow(
target="target",
source="source",
keys=["userId"],
stored_as_scd_type=2
)
Teszt:
Ez a teszt létrehoz egy kezdeti pillanatképet, futtatja a folyamatot, majd egy pillanatkép-frissítést szimulál az új adatok csonkolásával és beszúrásával annak ellenőrzéséhez, hogy a CDC rögzíti-e az összes módosítást.
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
def test_auto_cdc_from_snapshot_flow(test_spark):
# Create initial snapshot
test_spark.sql("""
CREATE TABLE catalog.schema.snapshot AS
SELECT * FROM VALUES
(1, 'Alice', '2024-01-01'),
(2, 'Bob', '2024-01-02')
AS t(userId, name, created_at)
""")
# Run the pipeline
test_pipeline.run(test_spark, set(["catalog.schema.target"]))
# Simulate a new snapshot by truncating and inserting updated data
test_spark.sql("TRUNCATE TABLE catalog.schema.snapshot")
test_spark.sql("INSERT INTO catalog.schema.snapshot VALUES (2, 'Bob', '2024-01-03')")
test_pipeline.run(test_spark, set(["catalog.schema.target"]))
# Verify SCD Type 2: should have 3 rows (original Alice, original Bob, updated Bob)
result = test_spark.table("catalog.schema.target")
assert result.count() == 3
user_ids = [row["userId"] for row in result.collect()]
assert set(user_ids) == {1, 2}
4. példa: Csatlakozások és elvárások tesztelése
Cél: Ellenőrizze, hogy az illesztések megfelelően működnek-e, és az elvárások kiszűrik az érvénytelen adatokat.
Folyamatátalakítás:
Ez az átalakítás összekapcsolja az ingatlan képeit a szolgáltatásokkal, és egy ellenőrzést alkalmaz a 2024. január előtt feltöltött képek kiszűrésére.
from pyspark import pipelines as dp
@dp.table
@dp.expect_or_drop("uploaded after Jan 2024", "uploaded_at > '2024-01-01'")
def property_images_amenities_join():
return (
spark.read.table("catalog.schema.property_images")
.join(
spark.read.table("catalog.schema.property_amenities"),
on="property_id",
how="inner"
)
)
Tesztek:
Ezek a tesztek ellenőrzik, hogy az illesztés a megfelelő számú sort hozza-e létre, és hogy a várakozás sikeresen kiszűri-e az érvénytelen feltöltési dátumokkal rendelkező rekordokat.
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
# Mock property datasets
def mock_properties(session):
session.sql("""
CREATE TABLE catalog.schema.property_images AS
SELECT * FROM VALUES
(101, 'img1.jpg', '2024-02-01'),
(102, 'img2.jpg', '2024-01-15'),
(103, 'img3.jpg', '2024-12-20')
AS t(property_id, image_url, uploaded_at)
""")
session.sql("""
CREATE TABLE catalog.schema.property_amenities AS
SELECT * FROM VALUES
(101, 'wifi'),
(102, 'pool'),
(103, 'parking')
AS t(property_id, amenity)
""")
# Test 1: Join
def test_property_join(test_spark):
mock_properties(test_spark)
test_pipeline.run(test_spark, set(["catalog.schema.property_images_amenities_join"]))
result = test_spark.table("catalog.schema.property_images_amenities_join")
# Should have 3 rows after join
assert result.count() == 3
# Check all property_ids are present
property_ids = set(row["property_id"] for row in result.collect())
assert property_ids == {101, 102, 103}
# Test 2: Expectation
def test_property_expectation(test_spark):
mock_properties(test_spark)
# Add a row with uploaded_at before Jan 2024
test_spark.sql("""
INSERT INTO catalog.schema.property_images VALUES (104, 'img4.jpg', '2023-12-31')
""")
# Add a matching row in the amenities table for the join
test_spark.sql("""
INSERT INTO catalog.schema.property_amenities VALUES (104, 'gym')
""")
test_pipeline.run(test_spark, set(["catalog.schema.property_images_amenities_join"]))
result = test_spark.table("catalog.schema.property_images_amenities_join")
# Only property_ids with uploaded_at > '2024-01-01' should be present
valid_ids = set(row["property_id"] for row in result.collect())
assert 104 not in valid_ids
assert valid_ids == {101, 102, 103}