Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Cette fonctionnalité est en version bêta.
Pour obtenir des informations générales sur Python tests unitaires dans Databricks, consultez Python tests unitaires.
Les pipelines Lakeflow prennent en charge l’écriture de tests unitaires Python dans l’éditeur de pipelines Lakeflow web. Cela vous permet de valider Python ou la logique de transformation SQL à l’aide de données fictifs. Avec le cadre de test des pipelines, vous pouvez tester des cas limites, valider les API propriétaires des pipelines (Auto CDC, tables de streaming, contraintes, flux d’ajout) et itérer à l’aide d’entrées fictives pour les opérations sur les identifiants de table prises en charge. Passez en revue les limitations d’isolation avant d’exécuter des tests.
- Exécution de test isolée : l’infrastructure fournit une session SparkSession qui redirige les opérations de table vers un schéma de test temporaire dans le catalogue par défaut du pipeline. Vous pouvez donc simuler des données d’entrée et écrire des sorties de test sans affecter les tables de production. L’isolation s’applique aux opérations qui font référence à une table par nom ; voir Limitations.
- Étendue de test flexible : exécutez un sous-ensemble d’un pipeline (tables individuelles, chaînes de tables dépendantes ou pipelines entiers) sur le calcul du pipeline à l’aide du test SparkSession.
- Validation des résultats : vérifiez les résultats des tables de sortie isolées créées dans un test à l’aide d’assertions pytest standard.
Quand utiliser des tests unitaires
Les cas d’usage classiques sont les suivants :
- Validation de la nouvelle logique de transformation : testez que votre transformation produit le schéma attendu, le nombre de lignes, les agrégations et la logique métier avant de s’exécuter sur des données de production.
- Test des spécifications Auto CDC : vérifiez que vos définitions de flux Auto CDC traitent correctement les événements de changement, y compris les insertions, les mises à jour, les suppressions et les types SCD (dimensions à évolution lente), à l’aide de données fictives.
- Test des attentes et des règles de qualité des données : vérifiez que les attentes échouent quand elles doivent et passent quand les données sont valides.
- Test entre les tables dépendantes : chaînes de tests de transformations (par exemple, bronze, argent et or) pour valider que les données circulent correctement via votre graphique de pipeline.
Requirements
Autorisation
Ownerdu pipeline, plus les privilègesUSE CATALOGetCREATE SCHEMAsur le catalogue par défaut du pipeline. L’infrastructure a besoin de ces privilèges pour créer le schéma de test temporaire où les tests s’exécutent.Pour vérifier ou définir l’autorisation de pipeline, ouvrez le pipeline, puis cliquez sur Partager. Vous devez être le
Owner(IS OWNER) du pipeline ;CAN RUNetCAN MANAGEne suffisent pas pour exécuter les tests. Consultez Configurer les autorisations de pipeline.Pour vérifier ou définir les privilèges du catalogue, ouvrez le catalogue dans l’Explorateur de catalogues, sélectionnez l’onglet Autorisations , puis vérifiez que vous disposez
USE CATALOGetCREATE SCHEMA. Un propriétaire de catalogue, un administrateur de metastore ou un utilisateur disposant duMANAGEprivilège peut les accorder, notamment avec SQL :GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;Pour plus d’informations, consultez la référence des privilèges du catalogue Unity.
Le pipeline doit être configuré en mode déclenché (non continu).
Le pipeline doit se trouver sur le canal PREVIEW . Les tests unitaires sont en version bêta et sont disponibles uniquement en préversion.
Spark Connect n’est pas pris en charge.
Note
L’isolation des tests couvre les opérations sur une table qui font référence à une table par son nom. Les opérations qui contournent l’isolation peuvent se produire à la fois dans votre code de test et dans tout code de pipeline exécuté par les sorties que vous sélectionnez, y compris ses dépendances transitives. Un fichier de test qui semble sans danger peut malgré tout exécuter un pipeline qui lit ou écrit des données via un chemin d’accès ou un connecteur, et agit sur des données de production. Pour empêcher les tests d’affecter les données de production ou les métadonnées, suivez les règles suivantes :
- Référencez chaque table par nom (
catalog.schema.table) et fictivez toutes les entrées par nom. Ne pas lire ou écrire via le chemin d’accès (/Volumes/...,dbfs:/...,s3://...,abfss://...) et ne pas lire depuis des connecteurs tels que Kafka ou Auto Loader. Ils contournent l’isolation et agissent sur des systèmes de production réels. - N’exécutez pas de déclarations de gouvernance ou de propriété, telles que
GRANT,REVOKE,ALTER ... OWNER TO,SET/UNSET TAGS, ouCREATE/DROP POLICY. Elles s’exécutent sur la ressource sécurisable de production réelle. - Ne créez pas de catalogues ou de schémas (
CREATE CATALOG,CREATE SCHEMA). Ils pointent vers votre metastore Unity Catalog réel. - N’exécutez pas l’intégralité du pipeline si son graphe inclut des entrées basées sur des chemins, des connecteurs, des écritures impératives ou d’autres effets secondaires externes. Sélectionnez uniquement les sorties dont les dépendances utilisent des opérations sur les tables de catalogue prises en charge et ont été remplacées par des entrées simulées.
Pour plus d’informations, consultez les Limitations .
Limitations
Warning
Certaines opérations contournent l’isolation des tests et peuvent agir sur des données ou métadonnées de production réelles. Passez en revue les limitations suivantes avant d’exécuter des tests.
L’isolation des tests est uniquement basée sur le nom de la table
Ne pas lisez ni écrivez via le chemin d'accès ou le connecteur. L’isolation redirige uniquement les opérations qui référencent une table par nom (par exemple,
spark.read.table("catalog.schema.table")oudf.write.saveAsTable("catalog.schema.table")). Les opérations traitées par un chemin d’accès ou via un connecteur contournent l’isolation et agissent directement sur les systèmes de production réels :-
L’écriture via un chemin (par exemple,
df.write.save("/Volumes/..."), un chemindbfs:/, ou un chemin cloud ou d’emplacement externe tel ques3://...ouabfss://...) écrit dans le stockage de production réel et peut écraser les données de production. -
La lecture par chemin (par exemple,
spark.read.load(path)ouspark.read.format("delta").load(path)) renvoie des données réelles de production au lieu de votre objet simulé. -
Lire depuis un connecteur établit une connexion à la véritable source de production. Cela inclut Kafka (lit à partir des répartiteurs réels) et Auto Loader (
cloudFiles, qui lit à partir du chemin réel de stockage dans le cloud). Ni l’un ni l’autre n’est redirigé vers vos données simulées.
-
L’écriture via un chemin (par exemple,
N’utilisez pas la fonction table
event_log()dans un test unitaire de pipeline. En mode test,event_log()n’est pas redirigé vers le journal des événements de votre exécution de test. Elle peut renvoyer le journal des événements de production ou un journal des événements précédemment enregistré, de sorte que les assertions effectuées à son encontre risquent de lire des données de production. Utilisez plutôt leevent_log_table_namerenvoyé par l’exécution et interrogez-le à l’aide detest_spark.event_log_table_namepeut êtreNone(par exemple, si le nom de la table du journal des événements ne peut pas être résolu), vérifiez-le avant d’interroger :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)N’affirmez
status.is_successpas avant de lire le journal des événements si votre objectif est de diagnostiquer une mise à jour ayant échoué. Le journal des événements est souvent ce que vous inspectez pour comprendre pourquoi une mise à jour a échoué.
Opérations de gouvernance et de DDL
- Les modifications du catalogue, du schéma, des autorisations, du propriétaire, des balises et des stratégies ne sont pas prises en charge. Cela inclut ,
CREATE/DROP/ALTER CATALOG(y comprisCREATE/DROP/ALTER SCHEMASET MANAGED LOCATION),GRANT/REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSet .CREATE/DROP POLICYCertains formulaires SQL exécutés par le biaistest_sparksont rejetés comme défense en profondeur ; d’autres formulaires, ou les mêmes opérations appelées via des API directes, peuvent atteindre des objets de production réels. Ne vous fiez pas à ces gardes comme limite d’isolement. Conservez ces instructions hors de votre code de test et hors du code de pipeline exécuté par les sorties sélectionnées.
Limitations opérationnelles
- L’exécution simultanée n’est pas prise en charge : l’exécution d’un test et d’une mise à jour de pipeline en même temps n’est pas prise en charge, et le système ne l’empêche pas. Il n’y a aucune coordination entre les deux, de sorte que l’exécution simultanée peut faire face à des ressources, dégrader gravement les performances de votre mise à jour de production ou provoquer l’échec du test de démarrage. Ne démarrez pas de test pendant que le pipeline exécute une mise à jour (ou démarrez une mise à jour pendant l’exécution d’un test) ; attendez la fin de toute mise à jour en cours avant d’exécuter des tests.
-
Schémas temporaires après l’arrêt anormal : chaque exécution de test crée un schéma temporaire (nommé
redirecting_<id>) dans le catalogue par défaut du pipeline et le supprime automatiquement une fois l’exécution terminée. Si une exécution se termine anormalement (par exemple, si les ressources de calcul sont perdues en cours d’exécution), le schéma temporaire peut subsister et contenir les tables simulées et les tables de sortie de l’exécution. Elle n’affecte pas les données de production. Pour libérer de l’espace de stockage, supprimez manuellement tous les schémas restants dont les noms commencent parredirecting_dans le catalogue par défaut du pipeline. - Les exécutions de test consomment le calcul : les exécutions de test s’exécutent sur le calcul du pipeline et sont facturées comme mises à jour normales du pipeline. Il n’existe pas de comptabilisation distincte pour les exécutions de test.
-
L’actualisation complète n’est pas prise en charge : seule l’actualisation sélective est disponible.
test_pipeline.run()actualise les sorties que vous sélectionnez (ou toutes les sorties lorsque vous ne passez aucune sélection) ; l’actualisation complète et la sélection de l’actualisation complète ne sont pas implémentées.
Limitations de création et de fidélité
- Exécution de l’éditeur uniquement : les tests doivent être exécutés à partir de l’éditeur de pipelines Lakeflow web.
- Python tests uniquement : les tests doivent être écrits dans Python. Vous pouvez tester des pipelines SQL, mais les tests eux-mêmes doivent être écrits dans Python.
- Fidélité aux règles de gouvernance : les données fictives n’héritent pas des filtres de ligne ni des masques de colonne définis sur les tables de production qu’elles remplacent. Les résultats des tests reflètent les entrées fictifs exactement comme vous les fournissez et peuvent différer de la façon dont la même requête se comporte sur les données de production régies.
Étape 1 : Mettre à jour les paramètres du pipeline
Configurez le pipeline pour qu’il s’exécute sur le canal PREVIEW en mode déclenché.
- Dans l’interface utilisateur, ouvrez votre pipeline, puis cliquez sur Paramètres>avancés>channel>Preview
- Définissez mode Pipeline sur Déclenché (n’utilisez pas Continu).
Vous pouvez également modifier directement les paramètres de pipeline JSON :
"continuous": false,
"channel": "PREVIEW"
Étape 2 : Créer un fichier de test
Dans l’Éditeur de pipelines Lakeflow, cliquez sur le bouton + (ajouter) et sélectionnez Test. Cela crée un fichier de test (et le tests dossier, s’il n’existe pas déjà) qui n’est pas inclus dans votre code source de pipeline. Vous n’avez pas besoin de créer le tests dossier vous-même.
Étape 3 : Générer des tests
Genie Code peut générer une structure de test :
Dans le fichier de test, cliquez sur le bouton Générer des tests .
Vous pouvez également utiliser
/testsdans le mode agent de Genie Code.
Utilisez Genie Code pour générer du code passe-partout, puis personnalisez-le pour vos cas limites.
Vous pouvez également écrire le code de test vous-même. Ajoutez les importations suivantes en haut de chaque fichier de test :
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
Étape 4 : Exécuter des tests
Exécutez des tests à partir de l’éditeur de pipelines Lakeflow :
- Cliquez sur le bouton
(lecture) dans la marge à côté d’une fonction de test pour exécuter un test individuel.
- Cliquez sur Exécuter des tests dans le fichier en haut du fichier de test pour exécuter tous les tests de ce fichier.
Les résultats des tests (réussite ou échec) apparaissent dans le panneau inférieur de l’éditeur. Passez en revue les erreurs d’assertion pour déboguer les échecs.
Test d’API
| API | Description |
|---|---|
TestPipeline.active() |
Retourne un objet TestPipeline correspondant au pipeline en cours de modification dans l’éditeur Lakeflow Pipelines. Cet objet est une référence au pipeline, y compris son code source, ses configurations, son catalogue/schéma par défaut, etc. |
test_pipeline.run(test_spark, set([table_names])) |
Exécute de façon synchrone une mise à jour du pipeline, effectuant une actualisation sélective si les noms de table sont spécifiés. Revient une fois l’exécution du pipeline réussie ou terminée avec une exception. |
Fixture test_spark |
Crée un test SparkSession avec redirection de table de catalogue qui redirige automatiquement les lectures et les écritures de table qui référencent une table par nom (par exemple, spark.read.table("catalog.schema.table") ou df.write.saveAsTable("catalog.schema.table")) vers un schéma de test temporaire. La redirection s’applique uniquement aux opérations de table basées sur le nom ; elle ne couvre pas les lectures ou les écritures traitées par chemin d’accès ou par le biais d’un connecteur, qui agissent directement sur le système réel. Consultez Limitations. |
Créer des données fictives
Vous pouvez simuler des données d’entrée à l’aide de SQL ou 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")
Pour générer de plus grands volumes de données synthétiques réalistes, vous pouvez utiliser la bibliothèque Faker . Exécutez %pip install faker d’abord dans votre pipeline, puis créez un DataFrame à partir d’UDF basées sur 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")
Exécuter le pipeline ou des tables précises
# 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)
Exemples
Exemple 1 : Test des agrégations avec le nombre de lignes, le schéma et la gestion null
Objectif : Vérifier que l’agrégation des utilisateurs compte correctement les utilisateurs par type, gère les adresses e-mail nulles et produit le schéma attendu.
Transformations de pipeline :
Ces transformations créent un pipeline simple à deux tables : users sélectionne les données utilisateur et counts regroupe les utilisateurs par type et comptent le nombre total d’utilisateurs et d’e-mails valides.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, count, count_if
@dp.table
def users():
return (
spark.read.table("catalog.schema.wanderbricks_users")
.select("user_id", "email", "name", "user_type")
)
@dp.table
def counts():
return (
spark.read.table("catalog.schema.users")
.withColumn("valid_email", col("email").isNotNull())
.groupBy("user_type")
.agg(
count("user_id").alias("total_count"),
count_if("valid_email").alias("count_valid_emails")
)
)
Tests :
Ces tests vérifient le nombre de lignes, la structure du schéma, la gestion des valeurs nulles et la logique d’agrégation en créant des données utilisateur fictives avec des valeurs nulles intentionnelles et en exécutant le pipeline de manière isolée.
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)
Exemple 2 : Test d’Auto CDC
Objectif : Valider qu’Auto CDC traite correctement le flux de modifications avec des insertions et des mises à jour.
Transformation de pipeline :
Cette transformation configure Auto CDC à partir d’un flux de changements, qui lit les changements en continu et les applique à la table cible selon le type 1 de SCD (ne conserve que la version la plus récente).
from pyspark import pipelines as dp
from pyspark.sql.functions import col
@dp.view
def users():
return spark.readStream.table("catalog.schema.change_feed")
dp.create_streaming_table("target_autocdc")
dp.create_auto_cdc_flow(
target="target_autocdc",
source="users",
keys=["userId"],
sequence_by=col("ts"),
stored_as_scd_type=1
)
Tests :
Le premier test crée un flux de modification fictif avec plusieurs enregistrements pour le même userId (simulant une mise à jour) et vérifie que seul l’enregistrement le plus récent est conservé dans la cible. Le deuxième test simule des événements arrivant en retard et dans le désordre en exécutant le pipeline, en ajoutant d’autres événements au flux de modifications et en réexécutant le pipeline.
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"
Exemple 3 : Test du CDC automatique à partir d’un instantané
Objectif : Valider que CDC traite correctement les modifications de l’instantané, y compris les insertions, les mises à jour et les suppressions.
Transformation de pipeline :
Cette transformation configure Auto CDC depuis un instantané, qui lit les données d’une table d’instantané et suit les modifications dans le temps en SCD de type 2 (conserve l’historique complet).
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 :
Ce test crée une capture instantanée initiale, exécute le pipeline, puis simule une mise à jour de la capture instantanée en tronquant les données et en en insérant de nouvelles afin de vérifier que CDC capture toutes les modifications.
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}
Exemple 4 : Tests de jointures et attentes
Objectif : Vérifiez que les jointures fonctionnent correctement et que les attentes filtrent les données non valides.
Transformation de pipeline :
Cette transformation associe les images de propriété aux équipements et applique une attente pour exclure les images téléversées avant janvier 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"
)
)
Tests :
Ces tests vérifient que la jointure produit le nombre correct de lignes et que l’attente filtre correctement les enregistrements avec des dates de chargement non valides.
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}