Eenheidstests voor pijplijnen

Important

Deze functie bevindt zich in de bètaversie.

Zie Python eenheidstests voor algemene informatie over het testen van Python eenheden in Databricks.

Lakeflow-pijplijnen ondersteunen het schrijven van Python eenheidstests in de webgebaseerde Lakeflow Pipelines Editor. Hiermee kunt u Python of SQL-transformatielogica valideren met behulp van mockgegevens. Met het framework voor pijplijntests kunt u edge-aanvragen testen, eigen pijplijn-API's (Auto CDC, streamingtabellen, verwachtingen, toevoegstromen) testen en herhalen met behulp van mock-invoer voor ondersteunde tabel-id-bewerkingen. Controleer de isolatiebeperkingen voordat u tests uitvoert.

  • Geïsoleerde testuitvoering: Het framework biedt een SparkSession waarmee tabelbewerkingen worden omgeleid naar een tijdelijk testschema in de standaardcatalogus van de pijplijn, zodat u invoergegevens kunt mocken en testuitvoer kunt schrijven zonder dat dit van invloed is op productietabellen. Isolatie is van toepassing op bewerkingen die verwijzen naar een tabel op naam; zie Beperkingen.
  • Flexibel testbereik: Voer een subset van een pijplijn (afzonderlijke tabellen, ketens van afhankelijke tabellen of volledige pijplijnen) uit op de berekening van de pijplijn met behulp van de SparkSession-test.
  • Resultaatvalidatie: Controleer de resultaten van geïsoleerde uitvoertabellen die in een test zijn gemaakt met behulp van standaard pytest-asserties.

Wanneer moet u eenheidstests gebruiken

Typische gebruiksvoorbeelden zijn onder andere:

  • Nieuwe transformatielogica valideren: test of uw transformatie het verwachte schema, het aantal rijen, aggregaties en de bedrijfslogica produceert voordat deze wordt uitgevoerd op basis van productiegegevens.
  • Auto CDC-specificaties testen: controleer of uw Auto CDC-flowdefinities wijzigingsgebeurtenissen correct verwerken, inclusief invoegingen, updates, verwijderingen en SCD-typen (Slowly Changing Dimension), met behulp van simulatiegegevens.
  • Verwachtingen en regels voor gegevenskwaliteit testen: Controleer of verwachtingen falen wanneer dat moet en slagen wanneer gegevens geldig zijn.
  • Testen in afhankelijke tabellen: Test ketens van transformaties (bijvoorbeeld brons, zilver en goud) om te controleren of de gegevens correct door uw pijplijngrafiek stromen.

Requirements

  • PijplijnmachtigingOwner, plus de USE CATALOG en CREATE SCHEMA bevoegdheden voor de standaardcatalogus van de pijplijn. Het framework heeft deze bevoegdheden nodig om het tijdelijke testschema te maken waarin tests worden uitgevoerd.

    Als u de pijplijnmachtiging wilt controleren of instellen, opent u de pijplijn en klikt u op Delen. Het moet de pijplijn Owner (IS OWNER) zijn; CAN RUN en CAN MANAGE zijn niet voldoende om tests uit te voeren. Zie Pijplijnmachtigingen configureren.

    Als u de catalogusbevoegdheden wilt controleren of instellen, opent u de catalogus in Catalog Explorer, selecteert u het tabblad Machtigingen en bevestigt u dat USE CATALOG en CREATE SCHEMA. Een cataloguseigenaar, een metastore-beheerder of een gebruiker met de MANAGE bevoegdheid kan deze verlenen, waaronder met SQL:

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

    Voor meer informatie zie de referentie voor Unity Catalog-bevoegdheden.

  • De pipeline moet zijn geconfigureerd in de triggergestuurde modus (niet-continu).

  • Pijplijn moet zich in het PREVIEW-kanaal bevinden. Eenheidstests zijn bèta en zijn alleen beschikbaar in PREVIEW.

  • Spark Connect wordt niet ondersteund.

Note

Testisolatie omvat tabelbewerkingen die naar een tabel verwijzen op basis van de naam. Bewerkingen die isolatie omzeilen, kunnen zowel in uw testcode als in elke pijplijncode worden uitgevoerd door de uitvoer die u selecteert, inclusief de transitieve afhankelijkheden. Een testbestand dat er veilig uitziet, kan nog steeds een pipelineflow uitvoeren die via een pad of connector gegevens leest of schrijft en daarbij bewerkingen uitvoert op productiegegevens. Volg deze regels om te voorkomen dat tests van invloed zijn op productiegegevens of metagegevens:

  • Verwijs naar elke tabel op naam (catalog.schema.table) en bekijk alle invoergegevens op naam. Lees of schrijf niet via een pad (/Volumes/..., dbfs:/..., s3://..., abfss://...) en lees ook niet uit connectors zoals Kafka of Auto Loader. Deze omzeilen de isolatie en werken op echte productiesystemen.
  • Voer geen governance- of eigendomsverklaringen uit, zoals GRANT, REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS, of CREATE/DROP POLICY. Deze worden uitgevoerd op het werkelijke productie-securable.
  • Maak geen catalogi of schema's (CREATE CATALOG, CREATE SCHEMA). Deze bereiken uw echte Unity Catalog-metastore.
  • Voer de hele pijplijn niet uit als de grafiek padgebaseerde invoer, connectors, imperatieve schrijfbewerkingen of andere externe bijwerkingen bevat. Selecteer alleen uitvoer waarvan afhankelijkheden ondersteunde catalogustabelbewerkingen gebruiken en zijn vervangen door mock-invoer.

Zie Beperkingen voor details.

Limitations

Warning

Sommige bewerkingen omzeilen testisolatie en kunnen reageren op echte productiegegevens of metagegevens. Bekijk de volgende beperkingen voordat u tests uitvoert.

De isolatie van tests gebeurt alleen op basis van de tabelnaam

  • Lees of schrijf niet per pad of connector. Isolatie stuurt alleen bewerkingen om die met de naam van een tabel verwijzen (bijvoorbeeld spark.read.table("catalog.schema.table") of df.write.saveAsTable("catalog.schema.table")). Bewerkingen die via een pad of connector worden benaderd, omzeilen de isolatie en werken rechtstreeks op echte productiesystemen:

    • Schrijven via een pad (bijvoorbeeld df.write.save("/Volumes/..."), een dbfs:/-pad, of een cloudpad of een pad naar een externe locatie zoals s3://... of abfss://...) schrijft rechtstreeks naar productieopslag en kan productiegegevens overschrijven.
    • Lezen via pad (bijvoorbeeld spark.read.load(path) of spark.read.format("delta").load(path)) geeft echte productiegegevens terug in plaats van je mock.
    • Lezen vanuit een connector maakt verbinding met de echte productiebron. Dit omvat Kafka (dat van de werkelijke brokers leest) en Auto Loader (cloudFiles, dat van het pad in de werkelijke cloudopslag leest). Geen van beide wordt omgeleid naar uw mockgegevens.
  • Gebruik de event_log() tabelwaardefunctie niet van een pijplijneenheidtest. In de testmodus wordt event_log() niet omgeleid naar het gebeurtenislogboek van uw testuitvoering. Het kan de productie of het eerder geregistreerde gebeurtenislogboek retourneren, zodat asserties tegen het productiegegevens kunnen worden gelezen. Gebruik in plaats daarvan de door de uitvoering geretourneerde event_log_table_name en bevraag deze via test_spark. event_log_table_name kan zijn None (bijvoorbeeld als de tabelnaam van het gebeurtenislogboek niet kan worden omgezet), dus controleer deze voordat u een query uitvoert:

    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)
    

    Bevestig niet status.is_success voordat u het gebeurtenislogboek leest als u een mislukte update wilt diagnosticeren. Het gebeurtenislogboek is vaak wat u inspecteert om te begrijpen waarom een update is mislukt.

Governance en DDL-bewerkingen

  • Catalogus, schema, machtiging, eigendom, tag en beleidsmutaties worden niet ondersteund. Dit omvat CREATE/DROP/ALTER CATALOG, CREATE/DROP/ALTER SCHEMA(inclusief SET MANAGED LOCATION), GRANT/REVOKE, , ALTER ... OWNER TO, ,SET/UNSET TAGS en .CREATE/DROP POLICY Sommige vormen van SQL-instructies die via test_spark worden uitgevoerd, worden uit voorzorg als extra verdedigingslaag afgewezen; andere vormen, of dezelfde bewerkingen die via rechtstreekse API’s worden aangeroepen, kunnen daadwerkelijk productieobjecten bereiken. Vertrouw niet op deze bewakers als isolatiegrens. Houd deze statements uit je testcode en uit alle pijplijncode die door de geselecteerde outputs wordt uitgevoerd.

Operationele beperkingen

  • Gelijktijdige uitvoering wordt niet ondersteund: het uitvoeren van een test en een pijplijnupdate op hetzelfde moment wordt niet ondersteund en het systeem voorkomt dit niet. Er is geen coördinatie tussen de twee, dus het gelijktijdig uitvoeren ervan kan strijden voor resources, waardoor de prestaties van uw productie-update ernstig afnemen of de test niet kan worden gestart. Start geen test terwijl de pijplijn een update uitvoert (of start een update terwijl een test wordt uitgevoerd); wacht totdat een actieve update is voltooid voordat u tests uitvoert.
  • Tijdelijke schema's na abnormale beëindiging: elke testuitvoering maakt een tijdelijk schema (benoemd redirecting_<id>) in de standaardcatalogus van de pijplijn en wordt automatisch verwijderd wanneer de uitvoering is voltooid. Als een uitvoering abnormaal eindigt (de berekening gaat bijvoorbeeld halverwege de uitvoering verloren), kan het tijdelijke schema achterblijven, waarbij de mock- en uitvoertabellen van de uitvoering worden bewaard. Dit heeft geen invloed op productiegegevens. Om opslagruimte vrij te maken, verwijder handmatig alle achtergebleven schema's waarvan de namen beginnen met redirecting_ in de standaardcatalogus van de pipeline.
  • Testuitvoeringen verbruiken rekenkracht: testuitvoeringen worden uitgevoerd op de berekening van de pijplijn en worden gefactureerd als normale pijplijnupdates. Er is geen afzonderlijke meting voor testuitvoeringen.
  • Volledig vernieuwen wordt niet ondersteund: alleen selectief vernieuwen is beschikbaar. test_pipeline.run() vernieuwt de uitvoer die u selecteert (of alle uitvoer wanneer u geen selectie doorgeeft); volledige vernieuwing en selectie voor volledig vernieuwen worden niet geïmplementeerd.

Beperkingen bij het maken en de nauwkeurigheid

  • Uitvoering alleen in de editor: tests moeten worden uitgevoerd in de webgebaseerde Lakeflow Pipelines-editor.
  • Python tests alleen: tests moeten worden geschreven in Python. U kunt SQL-pijplijnen testen, maar de tests zelf moeten worden geschreven in Python.
  • Governanceconsistentie: Mockgegevens erven geen rijfilters of kolommaskers die zijn gedefinieerd voor de productietabellen die ze vervangen. Testresultaten weerspiegelen de mock-invoer precies zoals u ze opgeeft en kan verschillen van hoe dezelfde query zich gedraagt op beheerde productiegegevens.

Stap 1: Pijplijninstellingen bijwerken

Configureer de pijplijn die moet worden uitgevoerd op het PREVIEW-kanaal in de geactiveerde modus.

  1. Open in de gebruikersinterface je pijplijn en klik op Instellingen>Geavanceerde instellingen>Kanaal>Voorbeeld
  2. Stel Pipeline mode in op Geactiveerd (gebruik niet Continue).

U kunt ook de JSON van de pijplijninstellingen rechtstreeks bewerken:

"continuous": false,
"channel": "PREVIEW"

Stap 2: Een testbestand maken

Klik in de Lakeflow Pipelines Editor op de + knop (toevoegen) en selecteer Testen. Hiermee maakt u een testbestand (en de tests map, als dit nog niet bestaat) die niet is opgenomen in de broncode van de pijplijn. U hoeft de tests map niet zelf te maken.

Menu Pijplijnassets toevoegen met de optie Testen om een pytest-bestand te maken.

Stap 3: Tests genereren

Genie Code kan test scaffolding genereren:

  • Klik in het testbestand op de knop Tests genereren .

    Leeg testbestand met de knop Tests genereren.

  • U kunt ook /tests gebruiken in de agentmodus van Genie Code.

    Testbestand ingevuld door Genie Code met op TestPipeline gebaseerde eenheidstests.

Gebruik Genie Code om standaardcode te genereren en pas deze vervolgens aan voor je randgevallen.

U kunt de testcode ook zelf schrijven. Voeg de volgende importbewerkingen toe aan het begin van elk testbestand:

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

test_pipeline = TestPipeline.active()

Stap 4: Tests uitvoeren

Voer tests uit vanuit de Lakeflow Pipelines-editor:

  • Klik op het pictogram Afspelen. (afspelen) in de rugmarge naast een testfunctie om een afzonderlijke test uit te voeren.
  • Klik boven aan het testbestand op Tests uitvoeren om alle tests in dat bestand uit te voeren.

Testresultaten (geslaagd of mislukt) worden weergegeven in het onderste deelvenster editor. Controleer assertiefouten om mislukkingen te debuggen.

API's testen

API Description
TestPipeline.active() Retourneert een TestPipeline object voor de pijplijn die momenteel wordt bewerkt in de Lakeflow Pipelines Editor. Dit object is een verwijzing naar de pijplijn, met inbegrip van de broncode, configuraties, standaardcatalogus/schema, enzovoort.
test_pipeline.run(test_spark, set([table_names])) Voert synchroon een update van de pijplijn uit, waarbij selectief wordt vernieuwd indien tabelnamen zijn opgegeven. Wordt geretourneerd nadat de uitvoering van de pijplijn succesvol is voltooid of wordt beëindigd met een uitzondering.
test_spark armatuur Hiermee maakt u een test-SparkSession met omleiding van catalogustabellen waarmee automatisch tabellees- en schrijfbewerkingen worden omgeleid die verwijzen naar een tabel op naam (bijvoorbeeld spark.read.table("catalog.schema.table") of df.write.saveAsTable("catalog.schema.table")) naar een tijdelijk testschema. Omleiding is alleen van toepassing op tabelbewerkingen op basis van naam; het behandelt geen lees- of schrijfbewerkingen die zijn geadresseerd via een pad of via een connector, die rechtstreeks op het echte systeem handelen. Zie Beperkingen.

Gesimuleerde gegevens maken

U kunt invoergegevens mocken met behulp van SQL of 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")

Om grotere hoeveelheden realistische synthetische gegevens te genereren, kunt u de Faker-bibliotheek gebruiken. Voer %pip install faker eerst uit in uw pijplijn en bouw vervolgens een DataFrame van door Faker ondersteunde UDF's:

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

De pijplijn of specifieke tabellen uitvoeren

# 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

Voorbeeld 1: Aggregaties testen met rijaantal, schema en null-verwerking

Doel: Valideer de gebruikersaggregatie correct op type, verwerkt null-e-mailberichten en produceert het verwachte schema.

Pijplijntransformaties:

Deze transformaties maken een eenvoudige pijplijn met twee tabellen: users selecteert gebruikersgegevens en counts groepeert gebruikers op type en telt het totale aantal gebruikers en geldige e-mailberichten.

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

Tests:

Met deze tests worden rijaantallen, schemastructuur, null-verwerking en aggregatielogica gevalideerd door gesimuleerde gebruikersgegevens te maken met opzettelijke nullen en de pijplijn geïsoleerd uit te voeren.

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)

Voorbeeld 2: Auto CDC testen

Doel: Controleer of automatisch CDC wijzigingenfeed correct verwerkt met invoegingen en updates.

Pijplijntransformatie:

Met deze transformatie wordt Auto CDC geconfigureerd op basis van een feed met wijzigingen, die streamende wijzigingen leest en deze op de doeltabel toepast als SCD van type 1 (waarbij alleen de meest recente versie behouden blijft).

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
)

Tests:

De eerste test maakt een mock-wijzigingsfeed met meerdere records voor hetzelfde userId (waarmee een update wordt gesimuleerd) en controleert of alleen het meest recente record in het doelsysteem wordt bewaard. De tweede test simuleert te laat arriverende en niet op volgorde binnenkomende gebeurtenissen door de pijplijn uit te voeren, meer gebeurtenissen aan de wijzigingenfeed toe te voegen en de pijplijn opnieuw uit te voeren.

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"

Voorbeeld 3: Auto CDC testen vanuit momentopname

Doel: controleer of cdc momentopnamewijzigingen correct verwerkt, waaronder invoegen, updates en verwijderingen.

Pijplijntransformatie:

Met deze transformatie wordt Auto CDC ingesteld vanuit een momentopname; deze leest uit een snapshot-tabel en houdt wijzigingen in de loop van de tijd bij als SCD-type 2 (waarbij de volledige geschiedenis behouden blijft).

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:

Met deze test maakt u een eerste momentopname, voert u de pijplijn uit en simuleert u vervolgens een momentopname-update door nieuwe gegevens af te kapen en in te voegen om te controleren of alle wijzigingen worden vastgelegd in CDC.

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}

Voorbeeld 4: Het testen van joins en verwachtingen

Doel: Controleer of joins correct werken en of verwachtingen ongeldige gegevens uitfilteren.

Pijplijntransformatie:

Deze transformatie combineert afbeeldingen van accommodaties met voorzieningen en past een controle toe om afbeeldingen te filteren die vóór januari 2024 zijn geüpload.

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

Tests:

Met deze tests wordt gecontroleerd of de join het juiste aantal rijen produceert en of de verwachting records met ongeldige uploaddatums heeft gefilterd.

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}