Testování jednotek pro kanály

Important

Tato funkce je v beta verzi.

Obecné informace o testování jednotek Python v Databricks najdete v tématu Python testování jednotek.

Lakeflow Pipelines podporuje psaní jednotkových testů v Pythonu ve webovém editoru Lakeflow Pipelines. To umožňuje ověřit logiku transformace Python nebo SQL pomocí napodobených dat. Pomocí rámce pro testování pipeline můžete testovat hraniční případy, ověřovat proprietární rozhraní API pipeline (Auto CDC, streamovací tabulky, datová očekávání, append toky) a provádět iterace pomocí simulovaných vstupů pro podporované operace nad identifikátory tabulek. Před spuštěním testů zkontrolujte omezení izolace.

  • Izolované spouštění testů: Rámec poskytuje SparkSession, která přesměrovává operace s tabulkami do dočasného testovacího schématu ve výchozím katalogu kanálu, takže můžete simulovat vstupní data a zapisovat testovací výstupy, aniž byste ovlivnili produkční tabulky. Izolace se vztahuje na operace, které odkazují na tabulku podle názvu; viz Omezení.
  • Flexibilní rozsah testování: Spusťte podmnožinu pipeline (jednotlivé tabulky, řetězce závislých tabulek nebo celé pipeline) na výpočetních prostředcích pipeline s využitím testovací relace SparkSession.
  • Ověření výsledku: Pomocí standardních kontrolních výrazů pytest ověřte výsledky izolovaných výstupních tabulek vytvořených v testu.

Kdy použít testování jednotek

Mezi obvyklé případy použití patří:

  • Ověřování nové logiky transformace: Před spuštěním s produkčními daty otestujte, že transformace vytvoří očekávané schéma, počty řádků, agregace a obchodní logiku.
  • Testování specifikací automatického CDC: Pomocí napodobených dat ověřte, jestli definice toku automatického CDC správně zpracovávají události změn, zpracovávají vložení, aktualizace, odstranění a typy SCD (Pomalu se měnící dimenze).
  • Testování očekávání a pravidel kvality dat: Ověřte, že očekávání selžou, když mají selhat, a projdou, když jsou data platná.
  • Testování napříč závislými tabulkami: Testovací řetězy transformací (například bronzová, stříbrná a zlatá) za účelem ověření správného toku dat přes graf kanálu.

Requirements

  • Oprávnění kanálu OwneraUSE CATALOGCREATE SCHEMA oprávnění k výchozímu katalogu kanálu. Architektura potřebuje tato oprávnění k vytvoření dočasného testovacího schématu, ve kterém se testy spouští.

    Pokud chcete zkontrolovat nebo nastavit oprávnění kanálu, otevřete kanál a klikněte na Sdílet. Musí jít o kanál Owner (IS OWNER); CAN RUN a CAN MANAGE nestačí pro spuštění testů. Viz Konfigurace oprávnění pipeline.

    Pokud chcete zkontrolovat nebo nastavit oprávnění katalogu, otevřete katalog v Průzkumníku katalogu, vyberte kartu Oprávnění a potvrďte, že máte USE CATALOG a CREATE SCHEMA. Vlastník katalogu, správce metastoru nebo uživatel s oprávněním MANAGE je může udělit, a to i pomocí SQL:

    GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;
    

    Další informace naleznete v tématu Referenční informace o oprávněních katalogu Unity.

  • Pipeline musí být nakonfigurována pro spouštěný (neprůběžný) režim.

  • Pipeline musí běžet na Databricks Runtime ve verzi 18.1 nebo vyšší. Starší runtime neobsahují modul pro jednotkové testování. Chcete-li zjistit, na jaké verzi modulu runtime byla aktualizace spuštěna, dotazujte protokol událostí kanálu. Viz informace o modulu runtime.

  • Spark Connect se nepodporuje.

Note

Izolace testů zahrnuje operace tabulek, které odkazují na tabulku podle názvu. K operacím, které obcházejí izolaci, může docházet jak ve vašem testovacím kódu, tak v jakémkoli kódu pipeline spouštěném výstupy, které vyberete, včetně jejích tranzitivních závislostí. Testovací soubor, který vypadá bezpečně, stále může spustit tok kanálu pipeline, který čte nebo zapisuje prostřednictvím cesty nebo konektoru a pracuje s produkčními daty. Aby testy neovlivnily produkční data nebo metadata, dodržujte tato pravidla:

  • Odkazujte na každou tabulku podle názvu (catalog.schema.table) a napodobujte všechny vstupy podle názvu. Nečtěte ani zapisujte podle cesty (/Volumes/..., dbfs:/..., s3://..., abfss://...) a nečtěte z konektorů, jako je Kafka nebo Auto Loader. Tyto obcházejí izolaci a zasahují do skutečných produkčních systémů.
  • Nespouštějte příkazy pro správu nebo vlastnictví, například GRANT, REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS nebo CREATE/DROP POLICY. Tyto kroky se provádějí proti skutečné výrobě zabezpečitelné.
  • Nevytvádřujte katalogy ani schémata (CREATE CATALOG, CREATE SCHEMA). Tyto odkazují na skutečný metastore služby Unity Catalog.
  • Nespouštějte celý pipeline, pokud jeho graf obsahuje vstupy založené na cestách, konektory, imperativní zápisy nebo jiné externí vedlejší účinky. Vyberte pouze výstupy, jejichž závislosti používají podporované operace katalog-tabulka a byly nahrazeny napodobenými vstupy.

Podrobnosti viz Omezení.

Omezení

Výstraha

Některé operace obcházejí izolaci testů a můžou pracovat s reálnými produkčními daty nebo metadaty. Před spuštěním testů zkontrolujte následující omezení.

Izolace testů je pouze podle názvu tabulky.

  • Nečte nebo zapisujte podle cesty nebo konektoru. Izolace přesměruje pouze operace odkazované na tabulku podle názvu (například spark.read.table("catalog.schema.table")df.write.saveAsTable("catalog.schema.table")). Operace prováděné prostřednictvím cesty nebo konektoru obcházejí izolaci a působí přímo na reálné produkční systémy:

    • Zápis pomocí cesty (například df.write.save("/Volumes/..."), cesta dbfs:/ nebo cloudová cesta či cesta k externímu umístění, například s3://... nebo abfss://...) zapisuje do skutečného produkčního úložiště a může přepsat produkční data.
    • Načítání podle cesty (například spark.read.load(path) nebo spark.read.format("delta").load(path)) vrací skutečná produkční data místo vašeho mock objektu.
    • Čtení z konektoru se připojuje k reálnému produkčnímu zdroji. To zahrnuje Kafka (která čte ze skutečných brokerů) a Auto Loader (cloudFiles, který čte ze skutečné cesty v cloudovém úložišti). Ani jeden z nich není přesměrován na vaše testovací data.
  • event_log() Nepoužívejte funkci s hodnotou tabulky z testu jednotek kanálu. V testovacím režimu není event_log() přesměrován do protokolu událostí vašeho testovacího běhu. Může vrátit produkční nebo dříve zaznamenaný protokol událostí, takže ověření vůči němu mohou číst produkční data. Místo toho použijte event_log_table_name vrácené spuštěním a dotazujte ho prostřednictvím test_spark. event_log_table_name může být None (například pokud se název tabulky protokolu událostí nedá přeložit), proto ho před dotazováním zkontrolujte:

    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)
    

    Pokud vaším cílem je diagnostikovat neúspěšnou aktualizaci, nevytvářejte status.is_success před čtením protokolu událostí. Protokol událostí často kontrolujete, abyste pochopili, proč aktualizace selhala.

Správa a operace DDL

  • Katalog, schéma, oprávnění, vlastnictví, označení a změny zásad nejsou podporovány. To zahrnuje CREATE/DROP/ALTER CATALOG, (CREATE/DROP/ALTER SCHEMAvčetně SET MANAGED LOCATION), GRANT/REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS, a .CREATE/DROP POLICY Některé formuláře SQL prováděné prostřednictvím test_spark jsou odmítnuty jako hloubková ochrana; jiné formuláře nebo stejné operace vyvolané prostřednictvím přímých rozhraní API mohou dosáhnout skutečných produkčních objektů. Nespoléhejte na tyto stráže jako na hranici izolace. Nepoužívejte tyto příkazy v testovacím kódu ani v žádném kódu pipeline spouštěném vybranými výstupy.

Provozní omezení

  • Souběžné spouštění není podporováno: Spuštění testu a aktualizace kanálu ve stejnou dobu není podporována a systém ji nezabrání. Mezi těmito dvěma není žádná koordinace, takže jejich souběžné spuštění může soupeřit o prostředky, výrazně snížit výkon vaší produkční aktualizace nebo způsobit selhání spuštění testu. Nespouštějte test, když v pipeline probíhá aktualizace (ani nespouštějte aktualizaci, když právě běží test); před spuštěním testů počkejte, až se dokončí jakákoli probíhající aktualizace.
  • Dočasná schémata po abnormálním ukončení: Každý testovací běh vytvoří dočasné schéma (s názvem redirecting_<id>) ve výchozím katalogu kanálu a po skončení běhu je automaticky odstraní. Pokud běh skončí abnormálně (například když se během běhu ztratí výpočetní prostředek), může dočasné schéma zůstat zachováno a obsahovat testovací a výstupní tabulky daného běhu. Nemá vliv na produkční data. Pokud chcete uvolnit úložiště, ručně odstraňte všechna zbývající schémata, jejichž názvy začínají redirecting_ ve výchozím katalogu kanálu.
  • Testovací běhy využívají výpočetní prostředky: Testovací běhy se spouštějí na výpočetních prostředcích kanálu a jsou účtovány jako běžné aktualizace kanálu. Pro testovací běhy neexistuje žádné samostatné měření.
  • Úplná aktualizace není podporována: K dispozici je pouze selektivní aktualizace. test_pipeline.run() obnoví vybrané výstupy (nebo všechny výstupy, pokud nezadáte žádný výběr); úplné obnovení a výběr pro úplné obnovení nejsou implementovány.

Omezení vytváření a věrnosti

  • Spouštění pouze v editoru: Testy je nutné spouštět z webového editoru Lakeflow Pipelines.
  • pouze Python testy: Testy musí být napsané v Python. Kanály SQL můžete testovat, ale samotné testy musí být napsané v Python.
  • Věrnost zásad správného řízení: Napodobená data nedědí filtry řádků ani masky sloupců definované v produkčních tabulkách, které nahrazuje. Výsledky testů odrážejí napodobené vstupy přesně tak, jak je zadáte, a můžou se lišit od toho, jak se stejný dotaz chová u řídicích produkčních dat.

Krok 1: Aktualizace nastavení kanálu

Nakonfigurujte pipeline tak, aby se spouštěla v aktivovaném režimu.

  1. V uživatelském rozhraní otevřete pipeline a klikněte na Nastavení.
  2. Nastavte režim pipeline na Spouštěný (nepoužívejte Průběžný).

Případně upravte přímo nastavení pipeline ve formátu JSON:

"continuous": false

Krok 2: Vytvoření testovacího souboru

V Editoru kanálů Lakeflow klikněte na + tlačítko (přidat) a vyberte Test. Tím se vytvoří testovací soubor (a složka tests, pokud ještě neexistuje), který není součástí zdrojového kódu kanálu pipeline. testsSložku nemusíte vytvářet sami.

Nabídka pro přidání prostředků kanálu pipeline, která zobrazuje možnost Test pro vytvoření souboru pytest.

Krok 3: Generování testů

Genie Code může vygenerovat kostru testů:

  • V testovacím souboru klikněte na tlačítko Generovat testy .

    Prázdný testovací soubor s tlačítkem Generovat testy

  • Alternativně použijte /tests v režimu agenta Genie Code.

    Testovací soubor naplněný kódem Genie s testy jednotek založenými na TestPipeline.

Použijte Genie Code k vygenerování šablonového kódu a poté jej upravte pro své okrajové případy.

Případně můžete testovací kód napsat sami. Na začátek každého testovacího souboru přidejte následující importy:

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark

test_pipeline = TestPipeline.active()

Krok 4: Spuštění testů

Spusťte testy z Editoru kanálů Lakeflow:

  • Klikněte na ikonu Přehrát. Tlačítko (přehrát) v hřbetu vedle testovací funkce pro spuštění jednotlivého testu.
  • Kliknutím na Spustit testy v souboru v horní části testovacího souboru spustíte všechny testy v daném souboru.

Výsledky testů (úspěch nebo selhání) se zobrazí na dolním panelu Editoru. Zkontrolujte selhání asercí, abyste mohli ladit selhání.

Testování rozhraní API

API Description
TestPipeline.active() Vrátí objekt TestPipeline kanálu, který se aktuálně upravuje v editoru kanálů Lakeflow. Tento objekt je odkazem na kanál, včetně jeho zdrojového kódu, konfigurací, výchozího katalogu nebo schématu atd.
test_pipeline.run(test_spark, set([table_names])) Synchronně spustí aktualizaci pipeline a pokud jsou zadány názvy tabulek, provede selektivní obnovení. Vrátí se po úspěšném spuštění kanálu nebo ukončení s výjimkou.
test_spark svítidlo Vytvoří testovací SparkSession s přesměrováním odkazů na katalogové tabulky, které automaticky přesměrovává čtení a zápisy tabulek odkazujících na tabulku podle názvu (například spark.read.table("catalog.schema.table") nebo df.write.saveAsTable("catalog.schema.table")) do dočasného testovacího schématu. Přesměrování se vztahuje pouze na operace s tabulkami založené na názvu; nevztahuje se na čtení nebo zápisy adresované pomocí cesty nebo prostřednictvím konektoru, které působí přímo na skutečný systém. Viz Omezení.

Vytvoření napodobených dat

Vstupní data můžete napodobenit pomocí jazyka SQL nebo createDataFrame:

# 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")

K vygenerování větších objemů realistických syntetických dat můžete použít knihovnu Faker . Nejprve ve své pipeline spusťte %pip install faker, potom vytvořte DataFrame z UDF založených na nástroji Faker:

# 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")

Spusťte pipeline nebo konkrétní tabulky

# 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)

Příklady

Příklad 1: Testování agregací s počtem řádků, schématem a zpracováním hodnot null

Cíl: Ověřit, že agregace uživatelů správně počítá uživatele podle typu, zpracovává e-mailové adresy s hodnotou null a vytváří očekávané schéma.

Transformace kanálu:

Tyto transformace vytvoří jednoduchý kanál se dvěma tabulkami: users vybere uživatelská data a counts seskupí uživatele podle typu a spočítá celkový počet uživatelů a platných e-mailů.

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")
        )
    )

Testy:

Tyto testy ověřují počty řádků, strukturu schématu, zpracování hodnot null a logiku agregace vytvořením napodobených uživatelských dat s úmyslnými null a spuštěním kanálu v izolaci.

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)

Příklad 2: Testování funkce Auto CDC

Cíl: Ověřit, že Auto CDC správně zpracovává kanál změn obsahující vložení a aktualizace.

Transformace kanálu:

Tato transformace nastaví Auto CDC ze zdroje změn, který čte průběžné změny a aplikuje je na cílovou tabulku jako SCD typu 1 (zachovává pouze nejnovější verzi).

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
)

Testy:

První test vytvoří simulovaný kanál změn obsahující více záznamů pro stejný userId (který simuluje aktualizaci) a ověří, že v cíli zůstane zachován pouze nejnovější záznam. Druhý test simuluje pozdě příchozí události a události přicházející mimo pořadí tím, že spustí pipeline, přidá další události do informačního kanálu změn a pipeline znovu spustí.

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"

Příklad 3: Testování Auto CDC ze snímku

Cíl: Ověřte, že CDC správně zpracovává změny snímků, včetně vkládání, aktualizací a odstraňování.

Transformace kanálu:

Tato transformace nastaví Auto CDC na základě snímku; čte z tabulky snímku a průběžně sleduje změny jako SCD typu 2 (uchovává úplnou historii).

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
)

Test:

Tento test vytvoří počáteční snapshot, spustí pipeline a poté simuluje aktualizaci snapshotu pomocí vyprázdnění a vložení nových dat, aby ověřil, že CDC zachytí všechny změny.

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}

Příklad 4: Testování spojení a očekávání

Cíl: Ověřte, že spojení fungují správně a očekávání odfiltrují neplatná data.

Transformace kanálu:

Tato transformace spojuje obrázky nemovitostí s informacemi o vybavení a aplikuje pravidlo pro odfiltrování obrázků nahraných před lednem 2024.

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"
        )
    )

Testy:

Tyto testy ověřují, že spojení vytvoří správný počet řádků a že očekávané výsledky úspěšně vyfiltrují záznamy s neplatnými daty nahrávání.

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}