Folyamatláncok egységtesztelése

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élyekOwner, valamint a USE CATALOGCREATE SCHEMA folyamat 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); a CAN RUN és CAN MANAGE nem 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 és CREATE 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 az MANAGE SQL-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 POLICY Ezek 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áulspark.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/..."), egy dbfs:/ elérési útra, vagy egy felhőbeli vagy külső helyre mutató elérési útra, például s3://... vagy abfss://...) 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.
  • 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 visszaadott event_log_table_name elemet, és a(z) test_spark használatával kérdezze le. event_log_table_name lehet None (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_success az 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 POLICY A test_spark elemen 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ődik redirecting_ .
  • 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.

  1. 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
  2. Á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.

A pipeline-összetevők hozzáadása menü, amely a Teszt lehetőséget jeleníti meg pytest-fájl létrehozásához.

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.

    Üres tesztfájl a Tesztek létrehozása gombbal.

  • Másik lehetőségként használja a /tests elemet a Genie Code ügynöki módjában.

    A Genie Code által kitöltött tesztfájl TestPipeline-alapú egységtesztekkel.

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 ikonra (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}