Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Important
Den här funktionen finns i Beta.
Allmän information om Python enhetstestning i Databricks finns i Python enhetstestning.
Lakeflow Pipelines stöder att skriva Python-enhetstester i den webbaserade redigeraren för Lakeflow Pipelines. På så sätt kan du verifiera Python- eller SQL-transformeringslogik med hjälp av falska data. Med ramverket för testning av pipelines kan du testa gränsfall, validera proprietära pipeline-API:er (Auto CDC, strömningstabeller, förväntningsregler, append-flöden) och iterera med simulerade indata för tabellidentifieraroperationer som stöds. Granska isoleringsbegränsningarna innan du kör tester.
- Isolerad testkörning: Ramverket tillhandahåller en SparkSession som omdirigerar tabellåtgärder till ett tillfälligt testschema i pipelinens standardkatalog, så att du kan simulera indata och skriva testutdata utan att påverka produktionstabeller. Isolering gäller för åtgärder som refererar till en tabell efter namn. se Begränsningar.
- Flexibelt testomfång: Kör en delmängd av en pipeline (enskilda tabeller, kedjor av beroende tabeller eller hela pipelines) på pipelinens beräkning med hjälp av testet SparkSession.
- Resultatverifiering: Verifiera resultatet av isolerade utdatatabeller som skapats i ett test med hjälp av standard pytest-försäkran.
När du ska använda enhetstestning
Vanliga användningsfall är:
- Validera ny transformeringslogik: Testa att omvandlingen genererar det förväntade schemat, antalet rader, aggregeringar och affärslogik innan du kör mot produktionsdata.
- Testning av Auto CDC-specifikationer: Kontrollera att dina Auto CDC-flödesdefinitioner bearbetar ändringshändelser korrekt genom att hantera infogningar, uppdateringar, raderingar och SCD-typer (Slowly Changing Dimension) med hjälp av simulerade data.
- Testning av förväntningar och regler för datakvalitet: Kontrollera att förväntningarna underkänns när de ska och godkänns när data är korrekta.
- Testning över beroende tabeller: Testa omvandlingskedjor (till exempel brons, silver och guld) för att verifiera att data flödar korrekt via pipelinediagrammet.
Requirements
behörighet för pipelinen
Owner, samt behörigheternaUSE CATALOGochCREATE SCHEMApå pipelinens standardkatalog. Ramverket behöver dessa behörigheter för att skapa det tillfälliga testschemat där testerna körs.Om du vill kontrollera eller ange pipelinebehörigheten öppnar du pipelinen och klickar på Dela. Du måste vara pipeline
Owner(IS OWNER);CAN RUNochCAN MANAGEär inte tillräckliga för att köra tester. Se Konfigurera pipelinebehörigheter.Om du vill kontrollera eller ange katalogprivilegier öppnar du katalogen i Katalogutforskaren, väljer fliken Behörigheter och bekräftar att du har
USE CATALOGochCREATE SCHEMA. En katalogägare, en metaarkivadministratör eller en användare med behörighetenMANAGEkan bevilja dem, inklusive med SQL:GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;Mer information finns i Referens för behörigheter för Unity Catalog.
Pipeline måste konfigureras i triggat (ej kontinuerligt) läge.
Pipeline måste vara på kanalen PREVIEW. Enhetstestning finns i Beta och är endast tillgängligt för förhandsversion.
Spark Connect stöds inte.
Note
Testisolering gäller tabelloperationer som refererar till en tabell efter namn. Åtgärder som kringgår isolering kan ske både i testkoden och i valfri pipelinekod som körs av de utdata du väljer, inklusive dess transitiva beroenden. En testfil som verkar säker kan fortfarande köra ett pipeline-flöde som läser eller skriver via sökväg eller anslutning och som arbetar mot produktionsdata. Följ dessa regler för att förhindra att tester påverkar produktionsdata eller metadata:
- Hänvisa till varje tabell med namn (
catalog.schema.table) och simulera alla indata med namn. Läs eller skriv inte via sökväg (/Volumes/...,dbfs:/...,s3://...,abfss://...) och läs inte från anslutningar som Kafka eller Auto Loader. Dessa kringgår isolering och fungerar på verkliga produktionssystem. - Kör inte styrnings- eller ägarskapsuttryck, till exempel
GRANT,REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSeller .CREATE/DROP POLICYDessa körs på det faktiska produktionsskyddsobjektet. - Skapa inte kataloger eller scheman (
CREATE CATALOG,CREATE SCHEMA). Dessa når ditt riktiga Unity Catalog-metaarkiv. - Kör inte hela pipelinen om dess graf innehåller sökvägsbaserade indata, anslutningar, imperativa skrivåtgärder eller andra externa biverkningar. Välj endast utdata vars beroenden använder åtgärder för katalogtabeller som stöds och har ersatts med simulerade indata.
Mer information finns i Begränsningar.
Limitations
Varning
Vissa åtgärder kringgår testisolering och kan fungera på verkliga produktionsdata eller metadata. Granska följande begränsningar innan du kör tester.
Testisolering är enbart baserad på tabellnamn
Läs eller skriv inte via sökväg eller anslutningsprogram. Isolering omdirigerar endast åtgärder som refererar till en tabell efter namn (till exempel
spark.read.table("catalog.schema.table")ellerdf.write.saveAsTable("catalog.schema.table")). Åtgärder som adresseras via en sökväg eller genom en anslutning kringgår isolering och påverkar faktiska produktionssystem direkt:-
Att skriva via sökväg (till exempel
df.write.save("/Volumes/..."), endbfs:/-sökväg eller en molnsökväg eller sökväg till en extern plats soms3://...ellerabfss://...) skriver till faktisk produktionslagring och kan skriva över produktionsdata. -
Läsning via sökväg (till exempel
spark.read.load(path)ellerspark.read.format("delta").load(path)) returnerar verkliga produktionsdata i stället för din mockdata. -
Läsning från en anslutning ansluter till den faktiska produktionskällan. Detta inkluderar Kafka (som läser från de faktiska brokerna) och Auto Loader (
cloudFiles, som läser från den faktiska sökvägen i molnlagringen). Ingen av dem omdirigeras till dina simulerade data.
-
Att skriva via sökväg (till exempel
Använd inte den tabellvärda funktionen
event_log()från ett enhetstest för pipeline. I testlägeevent_log()omdirigeras inte till testkörningens händelselogg. Den kan returnera produktionshändelseloggen eller en tidigare registrerad händelselogg, vilket innebär att asserteringar mot den kan läsa produktionsdata. Använd i stället denevent_log_table_namesom returneras av körningen och gör en fråga mot den viatest_spark.event_log_table_namekan varaNone(till exempel om tabellnamnet för händelseloggen inte kan matchas), så kontrollera det innan du frågar: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)Hävda inte
status.is_successinnan du läser händelseloggen om ditt mål är att diagnostisera en misslyckad uppdatering. Händelseloggen är ofta det du inspekterar för att förstå varför en uppdatering misslyckades.
Styrnings- och DDL-åtgärder
- Katalog-, schema-, behörighets-, ägarskaps-, tagg- och principmutationer stöds inte. Detta inkluderar
CREATE/DROP/ALTER CATALOG,CREATE/DROP/ALTER SCHEMA(inklusiveSET MANAGED LOCATION),GRANT/REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSoch .CREATE/DROP POLICYVissa SQL-formulär som körs viatest_sparkavvisas som skydd på djupet. Andra formulär, eller samma åtgärder som anropas via direkta API:er, kan nå verkliga produktionsobjekt. Förlita dig inte på dessa vakter som en isoleringsgräns. Håll dessa satser borta från testkoden och från all pipelinekod som körs av de valda utgångarna.
Driftbegränsningar
- Samtidig körning stöds inte: Körning av ett test och en pipelineuppdatering på samma gång stöds inte och systemet förhindrar det inte. Det finns ingen samordning mellan de två, så att köra dem samtidigt kan kämpa för resurser, kraftigt försämra prestandan för din produktionsuppdatering eller orsaka att testet inte startar. Starta inte ett test när pipelinen kör en uppdatering (eller starta en uppdatering medan ett test körs); vänta tills alla pågående uppdateringar har slutförts innan testerna körs.
-
Tillfälliga scheman efter onormal avslutning: Varje testkörning skapar ett tillfälligt schema (med namnet
redirecting_<id>) i pipelinens standardkatalog och släpper det automatiskt när körningen är klar. Om en körning slutar onormalt (till exempel om beräkningen går förlorad mitt i körningen) kan det tillfälliga schemat lämnas kvar och innehålla körningens mock- och utdatatabeller. Det påverkar inte produktionsdata. Frigör lagring genom att manuellt ta bort eventuella överblivna scheman vars namn börjar medredirecting_i pipelinens standardkatalog. - Testkörningar förbrukar beräkningsresurser: Testkörningar körs på pipelineens beräkningsresurser och debiteras som vanliga pipelineuppdateringar. Det finns ingen separat mätning för testkörningar.
-
Fullständig uppdatering stöds inte: Endast selektiv uppdatering är tillgänglig.
test_pipeline.run()uppdaterar de utdata du väljer (eller alla utdata när du inte anger något urval); fullständig uppdatering och urval för fullständig uppdatering är inte implementerade.
Begränsningar för redigering och noggrannhet
- Körning endast i redigeraren: Tester måste köras från den webbaserade redigeraren för Lakeflow Pipelines.
- endast Python tester: Testerna måste skrivas i Python. Du kan testa SQL-pipelines, men själva testerna måste skrivas i Python.
- Överensstämmelse med styrning: Simulerade data ärver inte radfilter eller kolumnmasker som definierats för de produktionstabeller de ersätter. Testresultaten återspeglar de falska indata exakt som du anger dem och kan skilja sig från hur samma fråga beter sig på reglerade produktionsdata.
Steg 1: Uppdatera pipelineinställningar
Konfigurera pipelinen så att den körs på förhandsgranskningskanalen i utlöst läge.
- Öppna pipelinen i användargränssnittet och klicka på Inställningar>Avancerade inställningar>Kanalförhandsgranskning>
- Ställ in Pipeline-läge på Utlöst (använd inte Kontinuerlig).
Du kan också redigera JSON-pipelineinställningarna direkt:
"continuous": false,
"channel": "PREVIEW"
Steg 2: Skapa en testfil
I Lakeflow Pipelines-redigeraren klickar du på + knappen (lägg till) och väljer Testa. Detta skapar en testfil (och mappen, om den tests inte redan finns) som inte ingår i pipelinens källkod. Du behöver inte skapa tests mappen själv.
Steg 3: Generera tester
Genie Code kan generera testställningar:
I testfilen klickar du på knappen Generera tester .
Du kan också använda
/testsi Genie Code-agentläge.
Använd Genie Code för att generera grundkod och anpassa den sedan för dina specialfall.
Du kan också skriva testkoden själv. Lägg till följande importer överst i varje testfil:
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
Steg 4: Kör tester
Kör tester från Lakeflow Pipelines-redigeraren:
- Klicka på
(play)-knappen i marginalen bredvid en testfunktion för att köra ett enskilt test.
- Klicka på Kör tester i filen överst i testfilen för att köra alla tester i filen.
Testresultat (lyckade eller misslyckade) visas i redigerarens nedre panel. Granska asserteringsfel för att felsöka problem.
Testa API:er
| API | Beskrivning |
|---|---|
TestPipeline.active() |
Returnerar ett TestPipeline objekt för pipelinen som för närvarande redigeras i Lakeflow Pipelines-redigeraren. Det här objektet är en referens till pipelinen, inklusive dess källkod, konfigurationer, standardkatalog/schema osv. |
test_pipeline.run(test_spark, set([table_names])) |
Kör synkront en uppdatering av pipelinen och utför en selektiv uppdatering om tabellnamn anges. Returnerar när körningen av pipelinen lyckas eller avslutas på grund av ett undantag. |
test_spark Fixtur |
Skapar ett test av SparkSession med omdirigering av katalogtabeller som automatiskt omdirigerar tabellläsningar och skrivningar som refererar till en tabell med namn (till exempel spark.read.table("catalog.schema.table") eller df.write.saveAsTable("catalog.schema.table")) till ett tillfälligt testschema. Omdirigering gäller endast för namnbaserade tabelloperationer; den omfattar inte läsningar eller skrivningar som adresseras med sökväg eller via en anslutning, vilka sker direkt i det verkliga systemet. Se Begränsningar. |
Skapa falska data
Du kan simulera indata med antingen SQL eller 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 du vill generera större volymer av realistiska syntetiska data kan du använda Faker-biblioteket . Kör %pip install faker först i pipelinen och skapa sedan en DataFrame från Faker-baserade UDF:er:
# 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")
Kör pipeline eller specifika tabeller
# 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)
Exempel
Exempel 1: Testa aggregeringar med radantal, schema och nullhantering
Mål: Verifiera att användaraggregering korrekt räknar användare efter typ, hanterar null-e-postmeddelanden och skapar det förväntade schemat.
Pipeline-transformeringar:
Dessa transformeringar skapar en enkel pipeline med två tabeller: users väljer användardata och counts grupperar användare efter typ och räknar totalt antal användare och giltiga e-postmeddelanden.
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")
)
)
Tester:
Dessa tester validerar radantal, schemastruktur, null-hantering och aggregeringslogik genom att skapa falska användardata med avsiktliga null-värden och köra pipelinen isolerat.
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)
Exempel 2: Testning av Auto CDC
Mål: Verifiera att Auto CDC bearbetar ändringsflödet korrekt med infogningar och uppdateringar.
Pipelineomvandling:
Den här omvandlingen konfigurerar Auto CDC från ett ändringsflöde, som läser strömmande ändringar och tillämpar dem på måltabellen som SCD-typ 1 (behåller endast den senaste versionen).
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
)
Tester:
Det första testet skapar en simulerad ändringsfeed med flera poster för samma userId för att simulera en uppdatering och verifierar att endast den senaste posten behålls i målet. Det andra testet simulerar händelser som inkommer sent och inte är i ordning genom att köra pipelinen, lägga till fler händelser i ändringsflödet och köra pipelinen igen.
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"
Exempel 3: Testa Auto CDC från ögonblicksbild
Mål: Kontrollera att CDC bearbetar ändringar av ögonblicksbilder korrekt, inklusive infogningar, uppdateringar och borttagningar.
Pipelineomvandling:
Den här omvandlingen ställer in Auto CDC från en ögonblicksbild, som läser från en snapshot-tabell och spårar ändringar över tid som SCD typ 2 (bevarar fullständig historik).
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:
Det här testet skapar en första ögonblicksbild, kör pipelinen och simulerar sedan en uppdatering av ögonblicksbilden genom att trunkera och infoga nya data för att verifiera att CDC samlar in alla ändringar.
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}
Exempel 4: Testa kopplingar och förväntningar
Mål: Verifiera att kopplingar fungerar korrekt och förväntningar filtrerar bort ogiltiga data.
Pipelineomvandling:
Den här omvandlingen kopplar samman egenskapsbilder med bekvämligheter och tillämpar en förväntan på att filtrera bort bilder som laddats upp före januari 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"
)
)
Tester:
Dessa tester verifierar att sammanfogningen ger rätt antal rader och att kontrollen på ett korrekt sätt filtrerar bort poster med ogiltiga uppladdningsdatum.
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}