Aller au contenu principal

Recommandations d’attente et modèles avancés

Les modèles d'attentes avancés combinent les attentes sur plusieurs jeux de données pour garantir la qualité des données à grande échelle. Ces modèles supposent que vous comprenez la syntaxe et la sémantique des vues matérialisées, des tables de streaming et des attentes.

Pour un aperçu de base du comportement et de la syntaxe des attentes, consultez Gérer la qualité des données avec les attentes de pipeline.

Attentes portables et réutilisables

Databricks recommande les bonnes pratiques suivantes lors de la mise en œuvre des attentes pour améliorer la portabilité et réduire la charge de maintenance :

Recommandation

Impact

Stockez les définitions d'attente séparément de la logique de pipeline.

Appliquez facilement des attentes à plusieurs jeux de données ou pipelines. Mettez à jour, auditez et maintenez les attentes sans modifier le code source du pipeline.

Ajoutez des tags personnalisés pour créer des groupes d'attentes connexes.

Filtrez les attentes en fonction des tags.

Appliquez les attentes de manière cohérente sur des datasets similaires.

Utilisez les mêmes attentes sur plusieurs datasets et pipelines pour évaluer une logique identique.

Recommandation

Impact

Stockez les définitions d'attente séparément de la logique de pipeline.

Appliquez facilement des attentes à plusieurs jeux de données ou pipelines. Mettez à jour, auditez et maintenez les attentes sans modifier le code source du pipeline.

Ajoutez des tags personnalisés pour créer des groupes d'attentes connexes.

Filtrez les attentes en fonction des tags.

Appliquez les attentes de manière cohérente sur des datasets similaires.

Utilisez les mêmes attentes sur plusieurs datasets et pipelines pour évaluer une logique identique.

Les exemples suivants illustrent l'utilisation d'une table Delta ou d'un dictionnaire pour créer un repository d'attentes central. Des fonctions Python personnalisées appliquent ensuite ces attentes aux datasets dans un pipeline d'exemple :

remarque

Le chargement dynamique des attentes à partir d'un fichier n'est pas pris en charge dans SQL.

L’exemple suivant crée une table nommée rules pour maintenir les règles :

SQL
CREATE OR REPLACE TABLE
rules
AS SELECT
col1 AS name,
col2 AS constraint,
col3 AS tag
FROM (
VALUES
("website_not_null","Website IS NOT NULL","validity"),
("fresh_data","to_date(updateTime,'M/d/yyyy h:m:s a') > '2010-01-01'","maintained"),
("social_media_access","NOT(Facebook IS NULL AND Twitter IS NULL AND Youtube IS NULL)","maintained")
)

L'exemple Python suivant définit les attentes en matière de qualité des données en fonction des règles du tableau rules. La fonction get_rules() lit les règles de la table rules et renvoie un dictionnaire Python contenant les règles correspondant à l'argument tag passé à la fonction.

Dans cet exemple, le dictionnaire est appliqué à l'aide de @dp.expect_all_or_drop() décorateurs pour appliquer des contraintes de qualité des données.

Par exemple, tous les enregistrements qui ne respectent pas les règles étiquetées avec validity sont supprimés de la table raw_farmers_market :

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import expr, col

def get_rules(tag):
"""
loads data quality rules from a table
:param tag: tag to match
:return: dictionary of rules that matched the tag
"""
df = spark.read.table("rules").filter(col("tag") == tag).collect()
return {
row['name']: row['constraint']
for row in df
}

@dp.table
@dp.expect_all_or_drop(get_rules('validity'))
def raw_farmers_market():
return (
spark.read.format('csv').option("header", "true")
.load('/databricks-datasets/data.gov/farmers_markets_geographic_data/data-001/')
)

@dp.table
@dp.expect_all_or_drop(get_rules('maintained'))
def organic_farmers_market():
return (
spark.read.table("raw_farmers_market")
.filter(expr("Organic = 'Y'"))
)

Tables de validation et flux de contrôle du pipeline

Certains modèles de cette section, tels que la validation du nombre de lignes et l' unicité de la clé primaire, définissent un dataset distinct, une *table de validation*, qui vérifie une propriété sur d'autres tables et utilise expect_or_fail pour faire remonter les problèmes. Avant de vous fier à une table de validation pour contrôler un pipeline, comprenez ce que les attentes peuvent et ne peuvent pas contrôler :

  • **Les attentes renforcent la qualité des données, et non l'orchestration.** Dans un pipeline, les attentes déterminent quels enregistrements atteignent un dataset cible : warn conserve les enregistrements invalides et enregistre les métriques, drop les supprime, et fail arrête le flux incriminé. L'objectif est de s'assurer que seules les données nettoyées transitent, et non d'exécuter ou d'ignorer conditionnellement d'autres parties du pipeline.
  • Le comportement deexpect_or_fail dépend du mode d'exécution du pipeline. Dans un Trigger pipeline, une attente échouée entraîne l'échec et l'annulation de la mise à jour de ce seul flux ; les autres flux du même pipeline continuent de se mettre à jour indépendamment. Dans un pipeline continu, une attente échouée arrête le flux et tous les flux dépendants. Voir Échec sur les enregistrements non valides.
  • Une table de validation ne bloque pas ses tables en aval. La lecture d’une table de validation depuis un autre dataset ne bloque pas ce dataset en attendant le résultat de la validation ; ainsi, une validation échouée n’empêche pas les tables en aval de se mettre à jour.

Pour arrêter le traitement en aval lorsqu'une validation échoue, séparez la logique de validation et le travail en aval en pipelines distincts et orchestrez-les avec un job, en faisant en sorte que la tâche du pipeline en aval dépende de la tâche du pipeline de validation. Parce qu'une tâche de pipeline échoue lorsque sa mise à jour échoue, la tâche en aval ne s'exécute pas à moins que le pipeline de validation ne réussisse. Plus généralement, lorsque vous avez besoin d'une exécution conditionnelle ou de dépendances complexes, coordonnez plusieurs pipelines avec un job au lieu d'intégrer la logique dans un seul pipeline. Consultez Exécuter des pipelines dans un workflow.

Validation du nombre de lignes

L'exemple suivant valide l'égalité du nombre de lignes entre table_a et table_b pour vérifier qu'aucune donnée n'est perdue pendant les transformations :

Graphe de validation du nombre de lignes LFP avec utilisation des attentes

Python
@dp.materialized_view(
name="count_verification",
comment="Validates equal row counts between tables"
)
@dp.expect_or_fail("no_rows_dropped", "a_count == b_count")
def validate_row_counts():
return spark.sql("""
SELECT * FROM
(SELECT COUNT(*) AS a_count FROM table_a),
(SELECT COUNT(*) AS b_count FROM table_b)""")

Détection des enregistrements manquants

L'exemple suivant valide que tous les enregistrements attendus sont présents dans la table report :

Graphe de détection des lignes manquantes LFP avec utilisation des attentes

Python
@dp.materialized_view(
name="report_compare_tests",
comment="Validates no records are missing after joining"
)
@dp.expect_or_fail("no_missing_records", "r_key IS NOT NULL")
def validate_report_completeness():
return (
spark.read.table("validation_copy").alias("v")
.join(
spark.read.table("report").alias("r"),
on="key",
how="left_outer"
)
.select(
"v.*",
"r.key as r_key"
)
)

Unicité de la clé primaire

L'exemple suivant valide les contraintes de clé principale sur plusieurs tables :

Graphe d'unicité de la clé primaire LFP avec l'utilisation des attentes.

Python
@dp.materialized_view(
name="report_pk_tests",
comment="Validates primary key uniqueness"
)
@dp.expect_or_fail("unique_pk", "num_entries = 1")
def validate_pk_uniqueness():
return (
spark.read.table("report")
.groupBy("pk")
.count()
.withColumnRenamed("count", "num_entries")
)

Modèle d'évolution des schémas

L'exemple suivant montre comment gérer l'évolution du schéma pour des colonnes supplémentaires. Utilisez ce modèle lorsque vous migrez des sources de données ou gérez plusieurs versions de données en amont, garantissant la compatibilité ascendante tout en faisant respecter la qualité des données :

Validation de l'évolution des schémas LFP avec l'utilisation des attentes

Python
@dp.table
@dp.expect_all_or_fail({
"required_columns": "col1 IS NOT NULL AND col2 IS NOT NULL",
"valid_col3": "CASE WHEN col3 IS NOT NULL THEN col3 > 0 ELSE TRUE END"
})
def evolving_table():
# Legacy data (V1 schema)
legacy_data = spark.read.table("legacy_source")

# New data (V2 schema)
new_data = spark.read.table("new_source")

# Combine both sources
return legacy_data.unionByName(new_data, allowMissingColumns=True)

Modèle de validation basé sur la plage

L'exemple suivant illustre comment valider de nouveaux points de données par rapport à des plages statistiques historiques, contribuant à identifier les valeurs aberrantes et les anomalies dans votre flux de données :

Validation basée sur la plage LFP avec l'utilisation des attentes.

Python
@dp.view
def stats_validation_view():
# Calculate statistical bounds from historical data
bounds = spark.sql("""
SELECT
avg(amount) - 3 * stddev(amount) as lower_bound,
avg(amount) + 3 * stddev(amount) as upper_bound
FROM historical_stats
WHERE
date >= CURRENT_DATE() - INTERVAL 30 DAYS
""")

# Join with new data and apply bounds
return spark.read.table("new_data").crossJoin(bounds)

@dp.table
@dp.expect_or_drop(
"within_statistical_range",
"amount BETWEEN lower_bound AND upper_bound"
)
def validated_amounts():
return spark.read.table("stats_validation_view")

Mettre en quarantaine les enregistrements non valides

Ce modèle combine les attentes avec des tables et des vues temporaires pour suivre les métriques de qualité des données pendant les mises à jour des pipelines et permettre des chemins de traitement séparés pour les enregistrements valides et non valides dans les opérations en aval.

Modèle de quarantaine de données LFP avec utilisation des attentes

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import expr

rules = {
"valid_pickup_zip": "(pickup_zip IS NOT NULL)",
"valid_dropoff_zip": "(dropoff_zip IS NOT NULL)",
}
quarantine_rules = "NOT({0})".format(" AND ".join(rules.values()))

@dp.view
def raw_trips_data():
return spark.readStream.table("samples.nyctaxi.trips")

@dp.table(
temporary=True,
partition_cols=["is_quarantined"],
)
@dp.expect_all(rules)
def trips_data_quarantine():
return (
spark.readStream.table("raw_trips_data").withColumn("is_quarantined", expr(quarantine_rules))
)

@dp.view
def valid_trips_data():
return spark.read.table("trips_data_quarantine").filter("is_quarantined=false")

@dp.view
def invalid_trips_data():
return spark.read.table("trips_data_quarantine").filter("is_quarantined=true")