Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
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
Owneroraz uprawnieniaUSE CATALOGiCREATE SCHEMAdo 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 RUNiCAN MANAGEnie 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 CATALOGiCREATE SCHEMA. Właściciel katalogu, administrator magazynu metadanych lub użytkownik z uprawnieniemMANAGEmoż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, lubCREATE/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")lubdf.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żkidbfs:/lub ścieżki w chmurze albo lokalizacji zewnętrznej, takiej jaks3://...lubabfss://...) zapisuje dane w rzeczywistej produkcyjnej przestrzeni dyskowej i może nadpisać dane produkcyjne. -
Odczyt po ścieżce (na przykład
spark.read.load(path)lubspark.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.
-
Zapisywanie przy użyciu ścieżki (na przykład
Nie używaj funkcji
event_log()zwracającej tabelę w teście jednostkowym potoku. W trybie testowymevent_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żyjevent_log_table_namezwróconego przez uruchomienie i odpytaj go za pomocątest_spark.event_log_table_namemoż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_successprzed 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 tymCREATE/DROP/ALTER SCHEMASET MANAGED LOCATION),GRANT/REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSi .CREATE/DROP POLICYNiektóre formularze SQL wykonywane za pośrednictwemtest_sparksą 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ę odredirecting_, 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.
- W interfejsie UI otwórz pipeline i kliknij Ustawienia.
- 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.
Krok 3. Generowanie testów
Kod Genie może generować szkielet testowy:
W pliku testowym kliknij przycisk Generuj testy .
Alternatywnie użyj
/testsw trybie agenta Genie Code.
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
(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}