Testowanie jednostkowe potoków

Ważna

Ta funkcja jest dostępna w wersji beta.

Aby uzyskać ogólne informacje na temat testowania jednostkowego Python w usłudze Databricks, zobacz Python testowanie jednostkowe.

Potoki Lakeflow umożliwiają tworzenie testów jednostkowych w Pythonie w internetowym Edytorze potoków Lakeflow. Dzięki temu można zweryfikować Python lub logikę przekształcania SQL przy użyciu pozornych danych. Za pomocą platformy testowania potoku można testować przypadki brzegowe, weryfikować zastrzeżone interfejsy API potoku (Auto CDC, tabele przesyłania strumieniowego, oczekiwania, dołączanie przepływów) i iterować przy użyciu pozornych danych wejściowych dla obsługiwanych operacji identyfikatora tabeli. Przed uruchomieniem testów zapoznaj się z ograniczeniami izolacji.

  • Izolowane wykonywanie testów: framework udostępnia obiekt SparkSession, który przekierowuje operacje na tabelach do tymczasowego schematu testowego w domyślnym katalogu potoku, dzięki czemu można symulować dane wejściowe i zapisywać dane wyjściowe testów bez wpływu na tabele produkcyjne. Izolacja ma zastosowanie do operacji odwołujących się do tabeli według nazwy; zobacz Ograniczenia.
  • Elastyczny zakres testowania: wykonaj podzbiór potoku (pojedyncze tabele, łańcuchy tabel zależnych lub całych potoków) na obliczeniach potoku przy użyciu testowej usługi SparkSession.
  • Weryfikacja wyniku: sprawdź wyniki izolowanych tabel wyjściowych utworzonych w teście przy użyciu standardowych asercji pytest.

Kiedy należy używać testów jednostkowych

Typowe przypadki użycia to:

  • Weryfikowanie nowej logiki przekształcania: przed uruchomieniem względem danych produkcyjnych przetestuj oczekiwany schemat, liczbę wierszy, agregacje i logikę biznesową.
  • Testowanie specyfikacji Auto CDC: Zweryfikuj, czy definicje przepływu Auto CDC prawidłowo przetwarzają zdarzenia zmian, obsługując operacje wstawiania, aktualizacji i usuwania oraz typy SCD (wolno zmieniające się wymiary) przy użyciu danych testowych.
  • Testowanie oczekiwań i reguł jakości danych: sprawdź, czy oczekiwania kończą się niepowodzeniem, gdy powinny i przechodzą, gdy dane są prawidłowe.
  • Testowanie między tabelami zależnymi: Testowanie łańcuchów przekształceń (na przykład brązowego, srebra i złota) w celu sprawdzenia, czy dane przepływa prawidłowo przez wykres potoku.

Requirements

  • Uprawnienie do potoku Owner oraz uprawnienia USE CATALOG i CREATE SCHEMA do domyślnego katalogu potoku. Struktura wymaga tych uprawnień, aby utworzyć tymczasowy schemat testu, w którym są uruchamiane testy.

    Aby sprawdzić lub ustawić uprawnienia potoku, otwórz potok i kliknij Udostępnij. Musi to być potok Owner (IS OWNER); CAN RUN i CAN MANAGE nie wystarczą, aby uruchomić testy. Zobacz Konfigurowanie uprawnień potoku.

    Aby sprawdzić lub ustawić uprawnienia wykazu, otwórz katalog w Eksploratorze wykazu, wybierz kartę Uprawnienia i potwierdź, że masz USE CATALOG i CREATE SCHEMA. Właściciel katalogu, administrator magazynu metadanych lub użytkownik z uprawnieniem MANAGE może je przyznać, w tym za pomocą języka SQL:

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

    Aby uzyskać więcej informacji, zobacz Przewodnik po uprawnieniach Unity Catalog.

  • Potok przetwarzania musi być skonfigurowany w trybie uruchamianym przez wyzwalacz (nieciągłym).

  • Potok musi działać w środowisku Databricks Runtime w wersji 18.1 lub nowszej. Wcześniejsze czasy uruchomieniowe nie obejmowały modułu testów jednostkowych. Aby sprawdzić, w jakiej wersji środowiska uruchomieniowego uruchomiono aktualizację, sprawdź dziennik zdarzeń potoku. Zobacz informacje o środowisku uruchomieniowym .

  • Połączenie Spark Connect nie jest obsługiwane.

Note

Izolacja testów obejmuje operacje na tabelach, które odwołują się do tabeli po nazwie. Operacje omijające izolację mogą wystąpić zarówno w kodzie testów, jak i w dowolnym kodzie potoku przetwarzania uruchamianym przez wybrane dane wyjściowe, wraz z jego zależnościami przechodnimi. Plik testowy, który wygląda na bezpieczny, może mimo to uruchomić przepływ w potoku, który odczytuje lub zapisuje dane za pomocą ścieżki lub łącznika i działa na danych produkcyjnych. Aby zachować wpływ testów na dane produkcyjne lub metadane, postępuj zgodnie z następującymi regułami:

  • Odwołuj się do każdej tabeli po nazwie (catalog.schema.table) i symuluj wszystkie wejścia po nazwie. Nie odczytuj ani nie zapisuj według ścieżki (/Volumes/..., dbfs:/..., s3://..., abfss://...) i nie odczytuj z łączników, takich jak Kafka lub Auto Loader. Omijają izolację i działają na rzeczywistych systemach produkcyjnych.
  • Nie stosuj oświadczeń dotyczących zarządzania lub własności, takich jak GRANT, REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS, lub CREATE/DROP POLICY. Są wykonywane na rzeczywistym zabezpieczanym obiekcie produkcyjnym.
  • Nie twórz katalogów ani schematów (CREATE CATALOG, CREATE SCHEMA). Trafiają one do rzeczywistego magazynu metadanych Unity Catalog.
  • Nie uruchamiaj całego pipeline’u, jeśli jego graf zawiera wejścia oparte na ścieżkach, łączniki, zapisy imperatywne lub inne zewnętrzne skutki uboczne. Wybierz tylko wyniki, których zależności używają obsługiwanych operacji na tabelach katalogu i zostały zastąpione symulowanymi danymi wejściowymi.

Zobacz Ograniczenia, aby poznać szczegóły.

Ograniczenia

Warning

Niektóre operacje pomijają izolację testów i mogą działać na rzeczywistych danych produkcyjnych lub metadanych. Przed uruchomieniem testów zapoznaj się z następującymi ograniczeniami.

Izolacja testów odbywa się wyłącznie na podstawie nazwy tabeli

  • Nie odczytuj ani nie zapisuj przy użyciu ścieżki lub konektora. Izolacja przekierowuje tylko operacje odwołujące się do tabeli według nazwy (na przykład spark.read.table("catalog.schema.table") lub df.write.saveAsTable("catalog.schema.table")). Operacje adresowane za pomocą ścieżki lub przez konektor obchodzą mechanizmy izolacji i działają bezpośrednio na realnych systemach produkcyjnych:

    • Zapisywanie przy użyciu ścieżki (na przykład df.write.save("/Volumes/..."), ścieżki dbfs:/ lub ścieżki w chmurze albo lokalizacji zewnętrznej, takiej jak s3://... lub abfss://...) zapisuje dane w rzeczywistej produkcyjnej przestrzeni dyskowej i może nadpisać dane produkcyjne.
    • Odczyt po ścieżce (na przykład spark.read.load(path) lub spark.read.format("delta").load(path)) zwraca rzeczywiste dane produkcyjne zamiast mocka.
    • Odczyt z łącznika nawiązuje połączenie z rzeczywistym źródłem produkcyjnym. Obejmuje to Kafka (odczytuje dane z rzeczywistych brokerów) i Auto Loader (cloudFiles, który odczytuje dane z rzeczywistej ścieżki w magazynie w chmurze). Żaden z nich nie jest przekierowywany do danych pozorowanych.
  • Nie używaj funkcji event_log() zwracającej tabelę w teście jednostkowym potoku. W trybie testowym event_log() nie jest przekierowywany do dziennika zdarzeń przebiegu testu. Może zwrócić dziennik zdarzeń z systemu produkcyjnego lub wcześniej zarejestrowany dziennik zdarzeń, więc asercje wykonywane na nim mogą odczytywać dane produkcyjne. Zamiast tego użyj event_log_table_name zwróconego przez uruchomienie i odpytaj go za pomocą test_spark. event_log_table_name może to być None (na przykład jeśli nie można rozpoznać nazwy tabeli dziennika zdarzeń), więc sprawdź ją przed wykonaniem zapytania:

    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)
    

    Nie potwierdzaj status.is_success przed odczytaniem dziennika zdarzeń, jeśli twoim celem jest zdiagnozowanie nieudanej aktualizacji. Dziennik zdarzeń jest często sprawdzany, aby zrozumieć, dlaczego aktualizacja nie powiodła się.

Nadzór i operacje DDL

  • Mutacje katalogu, schematu, uprawnień, własności, tagów i zasad nie są obsługiwane. Obejmuje to : ,CREATE/DROP/ALTER CATALOG(w tym CREATE/DROP/ALTER SCHEMASET MANAGED LOCATION),GRANT/REVOKE , ALTER ... OWNER TO,SET/UNSET TAGS i .CREATE/DROP POLICY Niektóre formularze SQL wykonywane za pośrednictwem test_spark są odrzucane jako ochrona w głębi systemu; inne formularze lub te same operacje wywoływane za pośrednictwem bezpośrednich interfejsów API mogą uzyskiwać dostęp do rzeczywistych obiektów produkcyjnych. Nie należy polegać na tych strażnikach jako granicy izolacji. Nie umieszczaj tych instrukcji ani w kodzie testowym, ani w żadnym kodzie pipeline’u uruchamianym przez wybrane wyjścia.

Ograniczenia operacyjne

  • Współbieżne wykonywanie nie jest obsługiwane: uruchamianie testu i aktualizacji potoku jednocześnie nie jest obsługiwane, a system temu nie zapobiega. Nie ma koordynacji między nimi, więc uruchomienie ich współbieżnie może konkurować o zasoby, poważnie obniżając wydajność aktualizacji produkcyjnej lub powodując niepowodzenie testu. Nie uruchamiaj testu, gdy potok jest aktualizowany (ani nie rozpoczynaj aktualizacji, gdy trwa test); przed uruchomieniem testów poczekaj na zakończenie trwającej aktualizacji.
  • Schematy tymczasowe po nieprawidłowym zakończeniu: Każde uruchomienie testu tworzy schemat tymczasowy (o nazwie redirecting_<id>) w domyślnym katalogu potoku i automatycznie usuwa go po zakończeniu uruchomienia. Jeśli uruchomienie zakończy się nieprawidłowo (na przykład gdy w jego trakcie zostaną utracone zasoby obliczeniowe), schemat tymczasowy może pozostać, zawierając tabele testowe i wyjściowe tego uruchomienia. Nie ma to wpływu na dane produkcyjne. Aby odzyskać miejsce na dane, ręcznie usuń wszystkie pozostałe schematy, których nazwy zaczynają się od redirecting_, w domyślnym katalogu potoku.
  • Uruchomienia testowe zużywają zasoby obliczeniowe: uruchomienia testowe są wykonywane na zasobach obliczeniowych potoku i są rozliczane jak standardowe aktualizacje potoku. Nie ma oddzielnego pomiaru dla przebiegów testów.
  • Pełne odświeżanie nie jest obsługiwane: dostępne jest tylko selektywne odświeżanie. test_pipeline.run() odświeża wybrane dane wyjściowe (lub wszystkie dane wyjściowe, jeśli nie przekażesz żadnego wyboru); pełne odświeżenie i wybór pełnego odświeżenia nie są zaimplementowane.

Ograniczenia dotyczące tworzenia i dokładności odwzorowania

  • Uruchamianie wyłącznie w edytorze: testy muszą być uruchamiane w internetowym edytorze Lakeflow Pipelines Editor.
  • Tylko testy w Pythonie: Testy muszą być napisane w Pythonie. Potoki SQL można przetestować, ale same testy muszą być zapisywane w Python.
  • Zgodność z zasadami ładu danych: Dane testowe nie dziedziczą filtrów wierszy ani masek kolumn zdefiniowanych dla zastępowanych przez nie tabel produkcyjnych. Wyniki testów odzwierciedlają pozorne dane wejściowe dokładnie tak, jak je podajesz i mogą różnić się od tego, jak to samo zapytanie zachowuje się na zarządzanych danych produkcyjnych.

Krok 1. Aktualizowanie ustawień potoku

Skonfiguruj potok tak, aby działał w trybie uruchamianym przez wyzwalacz.

  1. W interfejsie UI otwórz pipeline i kliknij Ustawienia.
  2. Ustaw tryb potoku na wyzwolony (nie używaj trybu ciągłego).

Możesz też edytować ustawienia potoku bezpośrednio w formacie JSON:

"continuous": false

Krok 2. Tworzenie pliku testowego

W Edytorze potoków Lakeflow kliknij przycisk + (Dodaj) i wybierz opcję Test. Spowoduje to utworzenie pliku testowego (oraz folderu tests, jeśli jeszcze nie istnieje), który nie jest uwzględniony w kodzie źródłowym pipeline’u. Nie musisz samodzielnie tworzyć tests folderu.

Menu „Dodaj zasoby potoku” wyświetlające opcję „Test” umożliwiającą utworzenie pliku pytest.

Krok 3. Generowanie testów

Kod Genie może generować szkielet testowy:

  • W pliku testowym kliknij przycisk Generuj testy .

    Pusty plik testowy z przyciskiem Generuj testy.

  • Alternatywnie użyj /tests w trybie agenta Genie Code.

    Plik testowy wypełniony przez kod Genie za pomocą testów jednostkowych opartych na języku TestPipeline.

Użyj Genie Code, aby wygenerować kod szablonowy, a następnie dostosować go do swoich przypadków brzegowych.

Alternatywnie możesz samodzielnie napisać kod testowy. Dodaj następujące importy na początku każdego pliku testowego:

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

test_pipeline = TestPipeline.active()

Krok 4. Uruchamianie testów

Uruchamianie testów w Edytorze potoków Lakeflow:

  • Kliknij ikonę odtwarzania (przycisk Odtwórz) na marginesie obok funkcji testowej, aby uruchomić pojedynczy test.
  • Kliknij pozycję Uruchom testy w pliku w górnej części pliku testowego, aby uruchomić wszystkie testy w tym pliku.

Wyniki testów (powodzenie lub niepowodzenie) są wyświetlane w dolnym panelu Edytor. Przejrzyj błędy asercji, aby zdiagnozować niepowodzenia.

Testowanie interfejsów API

API Description
TestPipeline.active() Zwraca obiekt TestPipeline dla potoku aktualnie edytowanego w Edytorze potoków Lakeflow. Ten obiekt jest referencją do potoku, obejmującą jego kod źródłowy, konfiguracje, domyślny katalog/schemat itp.
test_pipeline.run(test_spark, set([table_names])) Wykonuje synchronicznie aktualizację potoku danych, przeprowadzając selektywne odświeżenie, jeśli określono nazwy tabel. Zwraca po pomyślnym zakończeniu wykonywania potoku lub gdy jego wykonywanie zakończy się wyjątkiem.
test_spark oprawa Tworzy testową sesję SparkSession z przekierowaniem katalog–tabela, które automatycznie przekierowuje operacje odczytu i zapisu tabel odwołujące się do tabeli po nazwie (na przykład spark.read.table("catalog.schema.table") lub df.write.saveAsTable("catalog.schema.table")) do tymczasowego schematu testowego. Przekierowanie ma zastosowanie tylko do operacji tabel opartych na nazwach; nie obejmuje odczytów ani zapisów adresowanych przez ścieżkę lub za pośrednictwem łącznika, który działa bezpośrednio w rzeczywistym systemie. Zobacz Ograniczenia.

Utwórz dane testowe

Dane wejściowe można symulować przy użyciu języka SQL lub 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")

Aby wygenerować większe ilości realistycznych danych syntetycznych, możesz użyć biblioteki Faker . Uruchom najpierw %pip install faker w potoku, a następnie utwórz obiekt DataFrame z funkcji UDF opartych na Fakerze:

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

Uruchom potok lub określone tabele

# 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

Przykład 1: Testowanie agregacji z uwzględnieniem liczby wierszy, schematu i wartości null

Cel: Zweryfikować, czy agregacja użytkowników poprawnie zlicza użytkowników według typu, obsługuje adresy e-mail o wartości null i generuje oczekiwany schemat.

Przekształcenia potoku:

Te przekształcenia tworzą prosty potok złożony z dwóch tabel: users wybiera dane użytkowników, a counts grupuje użytkowników według typu i zlicza całkowitą liczbę użytkowników oraz prawidłowe adresy 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:

Te testy weryfikują liczbę wierszy, strukturę schematu, obsługę wartości null i logikę agregacji, tworząc pozorowane dane użytkownika z celowymi wartościami null i uruchamiając potok w izolacji.

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)

Przykład 2. Testowanie automatycznej usługi CDC

Cel: Zweryfikuj, czy Auto CDC prawidłowo przetwarza strumień zmian zawierający operacje wstawiania i aktualizacji.

Przekształcanie potoku:

Ta transformacja konfiguruje funkcję Auto CDC na podstawie strumienia zmian, która odczytuje zmiany strumieniowe i stosuje je w tabeli docelowej jako SCD typu 1 (zachowuje tylko najnowszą wersję).

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:

Pierwszy test tworzy symulowany strumień zmian z wieloma rekordami dotyczącymi tego samego userId (symulując aktualizację) i sprawdza, czy w obiekcie docelowym zostaje zachowany tylko najnowszy rekord. Drugi test symuluje zdarzenia, które docierają z opóźnieniem i w niewłaściwej kolejności, uruchamiając potok, dodając kolejne zdarzenia do strumienia zmian i ponownie uruchamiając potok.

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"

Przykład 3: Testowanie automatycznego mechanizmu CDC z migawki

Cel: Zweryfikuj, czy CDC prawidłowo przetwarza zmiany w migawce, w tym wstawienia, aktualizacje i usunięcia.

Przekształcanie potoku:

Ta transformacja konfiguruje mechanizm Auto CDC na podstawie migawki, który odczytuje dane z tabeli migawek i śledzi zmiany w czasie jako SCD typu 2 (zachowuje pełną historię).

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:

Ten test tworzy początkową migawkę, uruchamia potok danych, a następnie symuluje aktualizację migawki przez wykonanie operacji TRUNCATE i wstawienie nowych danych, aby sprawdzić, czy CDC rejestruje wszystkie zmiany.

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}

Przykład 4: Testowanie złączeń i oczekiwań

Cel: Sprawdzić, czy złączenia działają prawidłowo i czy reguły oczekiwań odfiltrowują nieprawidłowe dane.

Przekształcanie potoku:

Ta transformacja łączy zdjęcia nieruchomości z udogodnieniami i stosuje regułę, aby odfiltrować zdjęcia przesłane przed styczniem 2024 r.

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:

Te testy sprawdzają, czy sprzężenia generują prawidłową liczbę wierszy i czy oczekiwanie pomyślnie filtruje rekordy z nieprawidłowymi datami przekazywania.

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}