Aller au contenu principal

Tests unitaires pour les pipelines

info

Bêta

Cette fonctionnalité est en Bêta.

Pour des informations générales sur les tests unitaires Python dans Databricks, consultez les tests unitaires Python.

Lakeflow pipelines prennent en charge l'écriture de tests unitaires Python dans le Lakeflow Pipelines Editor basé sur le web. Cela vous permet de valider la logique de transformation Python ou SQL en utilisant des données fictives. Avec le framework de test de pipeline, vous pouvez tester les cas extrêmes, valider les APIs propriétaires de pipeline (Auto CDC, tables de streaming, attentes, flux d'ajout) et itérer en utilisant des entrées fictives pour les opérations d'identification de table prises en charge. Veuillez examiner les limitations d'isolation avant d'exécuter les tests.

  • **Exécution de test isolée** : Le framework fournit une SparkSession qui redirige les Opérations de table vers un schéma de test temporaire dans le catalogue default du pipeline, afin que vous puissiez 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 référencent une table par son nom ; voir Limitations.
  • Portée de test flexible : Exécutez un sous-ensemble d’un pipeline (tables individuelles, chaînes de tables dépendantes ou pipelines entiers) sur le compute du pipeline à l’aide du test SparkSession.
  • Validation des résultats : Vérifiez les résultats des tables de sortie isolées créées lors d'un test à l'aide d'assertions pytest standard.

Quand utiliser le test unitaire

Les cas d'utilisation typiques 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 l'exécuter sur des données de production.
  • Test des spécifications Auto CDC : Validez que vos définitions de flux Auto CDC traitent correctement les événements de changement, gérant les insertions, les mises à jour, les suppressions et les types SCD (Slowly Changing Dimension), à l’aide de données factices.
  • Test des attentes et des règles de qualité des données : Vérifiez que les attentes échouent lorsqu'elles le doivent et réussissent lorsque les données sont valides.
  • **Testez les chaînes de tables dépendantes** : Testez les chaînes de transformations (par exemple, bronze, argent et Gold) pour valider que les données circulent correctement à travers le graphe de votre pipeline.

Exigences

  • Owner Autorisation de pipeline, ainsi que les privilègesUSE CATALOG et CREATE SCHEMA sur le catalogue par default du pipeline. Le framework a besoin de ces privilèges pour créer le schéma de test temporaire où les tests sont exécutés.

    Pour vérifier ou définir l’autorisation du pipeline, ouvrez le pipeline et cliquez sur **Partager**. Vous devez être le **pipeline** Owner IS OWNER() ; CAN RUN et CAN MANAGE ne sont pas suffisants pour exécuter des tests. Voir Configurer les autorisations du pipeline.

    Pour vérifier ou définir les privilèges du catalogue, ouvrez le catalogue dans l'Explorateur de catalogues, sélectionnez l' tab Permissions et confirmez que vous disposez de USE CATALOG et CREATE SCHEMA. Un propriétaire de catalogue, un administrateur de metastore ou un utilisateur disposant du privilège MANAGE peut les accorder, y compris avec SQL :

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

    Pour plus d'informations, consultez la référence des privilèges Unity Catalog.

  • Le pipeline doit être configuré en mode Trigger (non continu).

  • Le pipeline doit se trouver sur le Canal de distribution . Le test unitaire est en version bêta et n'est disponible qu'en PRÉVISUALISATION.

  • Spark Connect n'est pas pris en charge.

remarque

Écrivez un code de test qui reste isolé

L’isolation des tests couvre les opérations de table qui référencént une table par 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 sûr peut toujours exécuter un flux de pipeline qui lit ou écrit par chemin ou connecteur, ce qui agit sur les données de production. Pour éviter que les tests n’affectent les données ou les métadonnées de production, suivez ces règles :

  • Référencez chaque table par son nom (catalog.schema.table), et simulez toutes les entrées par leur nom. Ne lisez ni n'écrivez par chemin d'accès (/Volumes/..., dbfs:/..., s3://..., abfss://...) et ne lisez pas à partir de connecteurs tels que Kafka ou Auto Loader. Celles-ci contournent l'isolation et agissent sur de véritables systèmes de production.
  • N’exécutez pas d’instructions de gouvernance ou de propriété, telles que GRANT, REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS ou CREATE/DROP POLICY. Ils s’exécutent par rapport à la ressource sécurisable de production réelle.
  • Ne créez pas de catalogues ou de schémas (CREATE CATALOG, CREATE SCHEMA). Ceux-ci atteignent votre véritable métastore Unity Catalog.
  • Ne lancez 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 de catalogue-table prises en charge et qui ont été remplacées par des entrées factices.

Voir les limites pour plus de détails.

Limitations

attention

Certaines opérations contournent l'isolation des tests et peuvent agir sur des données ou métadonnées de production réelles. Veuillez examiner les limitations suivantes avant d'exécuter les tests.

L'isolation des tests est uniquement par nom de table

  • Ne lisez pas et n'écrivez pas par chemin ou par connecteur. L'isolation ne redirige que les opérations qui référencent une table par son nom (par exemple, spark.read.table("catalog.schema.table") ou df.write.saveAsTable("catalog.schema.table")). Les opérations traitées par un chemin ou via un connecteur contournent l'isolation et agissent directement sur les systèmes de production réels :

    • L'écriture par chemin (par exemple, df.write.save("/Volumes/..."), un chemin dbfs:/, ou un chemin cloud ou d'emplacement externe tel que s3://... ou abfss://...) écrit dans un stockage de production réel et peut écraser les données de production.
    • Lecture par chemin d'accès (par exemple, spark.read.load(path) ou spark.read.format("delta").load(path)) renvoie des données de production réelles au lieu de votre simulation.
    • La lecture à partir d'un connecteur se connecte à la source de production réelle. Cela inclut Kafka (qui lit à partir des brokers réels) et Auto Loader (cloudFiles, qui lit à partir du chemin de stockage cloud réel). Aucune n'est redirigée vers vos données fictives.
  • N'utilisez pas la event_log() fonction à valeur de table à partir d'un test unitaire de pipeline. En mode test, event_log() n'est pas redirigé vers le log des événements de votre exécution de test. Il peut renvoyer le log de production ou enregistré précédemment, de sorte que les assertions à son encontre pourraient lire les données de production. Utilisez plutôt le event_log_table_name renvoyé par l'exécution et effectuez une query sur celui-ci via test_spark. event_log_table_name peut être None (par exemple, si le nom de la table du journal d'Logs ne peut pas être résolu), alors vérifiez-le avant d'interroger :

    Python
    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 pas status.is_success avant de lire le log des événements si votre objectif est de diagnostiquer une mise à jour ayant échoué. Le log des événements est souvent ce que vous inspectez pour comprendre pourquoi une mise à jour a échoué.

Gouvernance et opérations DDL

  • Les modifications de catalogue, de schéma, d'autorisation, de propriété, d'étiquette et de politique ne sont pas prises en charge. Cela inclut CREATE/DROP/ALTER CATALOG, CREATE/DROP/ALTER SCHEMA (y compris SET MANAGED LOCATION), GRANT/REVOKE, ALTER ... OWNER TO, SET/UNSET TAGS et CREATE/DROP POLICY. Certaines formes SQL exécutées via test_spark sont rejetées comme défense en profondeur ; d'autres formes, ou les mêmes Opérations invoquées via des APIs directes, peuvent atteindre de vrais objets de production. Ne vous fiez pas à ces protections comme limite d'isolation. Gardez ces instructions en dehors de votre code de test et de tout code de pipeline exécuté par les sorties sélectionnées.

Limitations opérationnelles

  • L'exécution concurrente n'est pas prise en charge : L'exécution simultanée d'un test et d'une mise à jour de pipeline n'est pas prise en charge, et le système ne l'empêche pas. Il n’y a aucune coordination entre les deux, donc les exécuter simultanément peut créer une concurrence pour les ressources, dégradant sévèrement les performances de votre mise à jour de production ou entraînant l’échec du start du test. Ne start pas un test pendant que le pipeline exécute une mise à jour (ou ne start pas une mise à jour pendant qu’un test est en cours d’exécution) ; attendez la fin de toute mise à jour en cours avant d’exécuter des tests.
  • Schémas temporaires après une terminaison anormale : Chaque exécution de test crée un schéma temporaire (nommé redirecting_<id>) dans le catalogue default du pipeline et le supprime automatiquement une fois l'exécution terminée. Si une exécution se termine anormalement (par exemple, le compute est perdu au milieu de l'exécution), le schéma temporaire peut être laissé, contenant les tables factices et de sortie de l'exécution. Cela n'affecte pas les données de production. Pour récupérer de l'espace de stockage, supprimez manuellement tout schéma restant dont le nom commence par redirecting_ dans le catalogue default du pipeline.
  • Les exécutions de test consomment du compute : Les exécutions de test s'exécutent sur le compute du pipeline et sont facturées comme des mises à jour normales de pipeline. Il n'y a pas de mesure distincte pour les exécutions de test.
  • Full refresh is not supported : Seule la refresh sélective est disponible. test_pipeline.run() refresh les sorties que vous sélectionnez (ou toutes les sorties si vous ne faites aucune sélection) ; le full refresh et la sélection de full-refresh ne sont pas implémentées.

Limitations de création et de fidélité

  • Exécution exclusive de l’éditeur : Les tests doivent être exécutés à partir de l’Éditeur de LakeFlow Pipelines basé sur le web.
  • Tests Python uniquement : les tests doivent être écrits en Python. Vous pouvez tester les pipelines SQL, mais les tests eux-mêmes doivent être écrits en Python.
  • Fidélité de la gouvernance : Les données factices n'héritent pas des filtres de lignes ou des masques de colonnes définis sur les tables de production qu'elles remplacent. Les résultats des tests reflètent les entrées factices exactement telles que vous les fournissez et peuvent différer de la façon dont la même query se comporte sur des données de production gouvernées.

Étape 1 : mettre à jour les paramètres du pipeline

Configurer le pipeline pour qu'il s'exécute sur le canal de distribution « PREVIEW » en mode Trigger.

  1. Dans l'interface utilisateur, ouvrez votre pipeline et cliquez sur **Paramètres** > **Paramètres avancés** > **Canal de distribution** > **Aperçu**
  2. Définissez le mode pipeline sur Déclenché (ne pas utiliser Continu).

Alternativement, modifiez directement les paramètres du pipeline JSON :

JSON
"continuous": false,
"channel": "PREVIEW"

Étape 2 : Créer un fichier de test

Dans l'éditeur de LakeFlow Pipelines, cliquez sur le bouton + (ajouter) et sélectionnez Test . Cela crée un fichier de test (et le dossier tests, s'il n'existe pas déjà) qui n'est pas inclus dans le code source de votre pipeline. Vous n'avez pas besoin de créer le dossier tests vous-même.

Ajouter le menu d&#39;assets de pipeline affichant l&#39;option Test pour créer un fichier pytest.

Étape 3 : générer des tests

Genie Code peut générer un échafaudage de test :

  • Dans le fichier de test, cliquez sur le bouton **Générer des tests**.

    Fichier de test vide avec le bouton Générer des tests.

  • Alternativement, utilisez /tests en mode agent Genie Code.

    Fichier de test renseigné par Genie Code avec des tests unitaires basés sur TestPipeline.

Utilisez Genie Code pour générer du code générique, puis personnalisez-le pour vos cas limites.

Alternativement, vous pouvez écrire le code de test vous-même. Ajoutez les importations suivantes en haut de chaque fichier de test :

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

test_pipeline = TestPipeline.active()

Étape 4 : Exécuter les tests

Exécutez les tests à partir de l’éditeur de LakeFlow Pipelines :

  • Cliquez sur le bouton Icône de lecture. (lecture) dans la gouttière à côté d’une fonction de test pour exécuter un test individuel.
  • Cliquez sur Exécuter les tests dans le fichier en haut du fichier de test pour exécuter tous les tests de ce fichier.

Les résultats des tests (succès ou échec) apparaissent dans le panneau inférieur de l’éditeur. Passez en revue les erreurs d’assertion pour déboguer les défaillances.

Test des APIs

API

Description

TestPipeline.active()

Renvoie un objet TestPipeline pour le pipeline actuellement en cours de modification dans l'éditeur de LakeFlow Pipelines. Cet objet est une référence au pipeline, y compris son code source, ses configurations, son catalogue/schéma default, etc.

test_pipeline.run(test_spark, set([table_names]))

Exécute de manière synchrone une mise à jour du pipeline, en effectuant une refresh sélective si les noms de table sont spécifiés. Retourne après que l'exécution du pipeline a réussi ou qu'elle a été interrompue avec une exception.

test_spark fixture

Crée une SparkSession de test avec redirection catalogue-table qui redirige automatiquement les lectures et écritures de table qui référencent une table **par nom** (par exemple, spark.read.table("catalog.schema.table") df.write.saveAsTable("catalog.schema.table")ou) 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 adressées par chemin ou via un connecteur, qui agissent directement sur le système réel. Consultez les Limites.

API

Description

TestPipeline.active()

Renvoie un objet TestPipeline pour le pipeline actuellement en cours de modification dans l'éditeur de LakeFlow Pipelines. Cet objet est une référence au pipeline, y compris son code source, ses configurations, son catalogue/schéma default, etc.

test_pipeline.run(test_spark, set([table_names]))

Exécute de manière synchrone une mise à jour du pipeline, en effectuant une refresh sélective si les noms de table sont spécifiés. Retourne après que l'exécution du pipeline a réussi ou qu'elle a été interrompue avec une exception.

test_spark fixture

Crée une SparkSession de test avec redirection catalogue-table qui redirige automatiquement les lectures et écritures de table qui référencent une table **par nom** (par exemple, spark.read.table("catalog.schema.table") df.write.saveAsTable("catalog.schema.table")ou) 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 adressées par chemin ou via un connecteur, qui agissent directement sur le système réel. Consultez les Limites.

Créez des données fictives

Vous pouvez simuler des données d'entrée soit à l'aide de SQL, soit à l'createDataFrame:

Python
# 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 des volumes plus importants de données synthétiques réalistes, vous pouvez utiliser la bibliothèque Faker. Exécutez %pip install faker dans votre pipeline d'abord, puis créez un DataFrame à partir d'UDF basées sur Faker :

Python
# 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écutez le pipeline ou des tables spécifiques

Python
# 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 des valeurs nulles

Objectif : Validez que l’agrégation des utilisateurs compte correctement les utilisateurs par type, gère les e-mails nuls et produit le schéma attendu.

**Pipeline Transformations** :

Ces transformations créent un pipeline simple à deux tables : users sélectionne les données utilisateur, et counts regroupe les utilisateurs par type et compte le nombre total d'utilisateurs et d'e-mails valides.

Python
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 valident 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 factices avec des valeurs nulles intentionnelles et en exécutant le pipeline de manière isolée.

Python
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 que la CDC automatique traite correctement le flux de changements avec les insertions et les mises à jour.

Transformation de pipeline :

Cette transformation configure le CDC automatique à partir d'un flux de modifications, qui lit les modifications en streaming et les applique à la table cible en tant que SCD de type 1 (ne conserve que la dernière version).

Python
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 modifications simulé avec plusieurs enregistrements pour le même userId (simulant une mise à jour) et vérifie que seul le dernier enregistrement est conservé dans la cible. Le second test simule des événements arrivant en retard et désordonnés en exécutant le pipeline, en ajoutant d'autres événements au flux de modifications et en exécutant à nouveau le pipeline.

Python
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 de l’Auto CDC à partir d’un snapshot

Objectif : Valider que le CDC traite correctement les modifications d'instantané, y compris les insertions, les mises à jour et les suppressions.

Transformation de pipeline :

Cette Transformation configure Auto CDC à partir d’un instantané, qui lit à partir d’une table instantanée et suit les changements au fil du temps en tant que SCD de type 2 (maintient l’historique complet).

Python
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 un instantané initial, exécute le pipeline, puis simule une mise à jour d'instantané en tronquant et en insérant de nouvelles données pour vérifier que le CDC capture toutes les modifications.

Python
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 : tester les jointures et les attentes

Objectif : valider que les jointures fonctionnent correctement et que les attentes filtrent les données non valides.

Transformation de pipeline :

Cette transformation joint les images de propriété avec les commodités et applique une attente pour filtrer les images upload avant janvier 2024.

Python
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 avec succès les enregistrements avec des dates d'upload non valides.

Python
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}