Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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 SCHEMAoprá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 RUNaCAN MANAGEnestačí 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 CATALOGaCREATE SCHEMA. Vlastník katalogu, správce metastoru nebo uživatel s oprávněnímMANAGEje 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 TAGSneboCREATE/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/..."), cestadbfs:/nebo cloudová cesta či cesta k externímu umístění, napříklads3://...neboabfss://...) 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)nebospark.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.
-
Zápis pomocí cesty (například
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žijteevent_log_table_namevrácené spuštěním a dotazujte ho prostřednictvímtest_spark.event_log_table_namemůže býtNone(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_successpř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 POLICYNěkteré formuláře SQL prováděné prostřednictvímtest_sparkjsou 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.
- V uživatelském rozhraní otevřete pipeline a klikněte na Nastavení.
- 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.
Krok 3: Generování testů
Genie Code může vygenerovat kostru testů:
V testovacím souboru klikněte na tlačítko Generovat testy .
Alternativně použijte
/testsv režimu agenta Genie Code.
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
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}