Aller au contenu principal

Gérer la qualité des données avec les attentes de pipeline

Utilisez les attentes pour appliquer des contraintes de qualité qui valident les données à mesure qu’elles circulent dans les pipelines ETL. Les attentes offrent un meilleur aperçu des métriques de qualité des données et vous permettent de faire échouer les mises à jour ou d’ignorer des enregistrements lors de la détection d’enregistrements non valides.

Pour les cas d'utilisation avancés et les meilleures pratiques recommandées, consultez Recommandations en matière d'attentes et modèles avancés.

Graphe de flux des attentes de pipeline

Quelles sont les attentes ?

Les attentes sont des clauses facultatives dans les instructions de création de vues matérialisées de pipeline, de tables de streaming ou de vues qui appliquent des contrôles de qualité des données sur chaque enregistrement passant par une query. Les attentes utilisent des énoncés booléens SQL standard pour spécifier des contraintes. Vous pouvez combiner plusieurs attentes pour un seul dataset et définir des attentes pour toutes les déclarations de datasets dans un pipeline.

remarque

Vous pouvez également définir des attentes sur les tables de streaming et les vues matérialisées soutenues par un pipeline autonome créé dans Databricks SQL. Utilisez la clause CONSTRAINT expectation_name EXPECT (expectation_expr) dans CREATE STREAMING TABLE et CREATE MATERIALIZED VIEW.

Les sections suivantes présentent les trois composants d'une attente et fournissent des exemples de syntaxe.

Nom de l’attente

Chaque attente doit avoir un nom, qui est utilisé comme identifiant pour suivre et surveiller l'attente. Choisissez un nom qui communique les métriques en cours de validation. L’exemple suivant définit l’attente valid_customer_age pour confirmer que l’âge est compris entre 0 et 120 ans :

important

Un nom d'attente doit être unique pour un dataset donné. Vous pouvez réutiliser les attentes sur plusieurs datasets dans un pipeline. Voir attentes portables et réutilisables.

Python
@dp.table
@dp.expect("valid_customer_age", "age BETWEEN 0 AND 120")
def customers():
return spark.readStream.table("datasets.samples.raw_customers")

Contrainte à évaluer

La clause de contrainte est une instruction conditionnelle SQL qui doit être évaluée à vrai ou faux pour chaque enregistrement. La contrainte contient la logique réelle de ce qui est validé. Lorsqu'un enregistrement échoue à cette condition, l'attente est Trigger.

Les contraintes doivent utiliser une syntaxe SQL valide et ne peuvent pas contenir les éléments suivants :

  • Fonctions Python personnalisées
  • Appels de service externes
  • Sous-requêtes référençant d'autres tables

Voici des exemples de contraintes qui pourraient être ajoutées aux instructions de création de dataset :

La syntaxe d'une contrainte en Python est :

Python
@dp.expect(<constraint-name>, <constraint-clause>)

Plusieurs contraintes peuvent être spécifiées :

Python
@dp.expect(<constraint-name>, <constraint-clause>)
@dp.expect(<constraint2-name>, <constraint2-clause>)

Exemples :

Python
# Simple constraint
@dp.expect("non_negative_price", "price >= 0")

# SQL functions
@dp.expect("valid_date", "year(transaction_date) >= 2020")

# CASE statements
@dp.expect("valid_order_status", """
CASE
WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
ELSE false
END
""")

# Multiple constraints
@dp.expect("non_negative_price", "price >= 0")
@dp.expect("valid_purchase_date", "date <= current_date()")

# Complex business logic
@dp.expect(
"valid_subscription_dates",
"""start_date <= end_date
AND end_date <= current_date()
AND start_date >= '2020-01-01'"""
)

# Complex boolean logic
@dp.expect("valid_order_state", """
(status = 'ACTIVE' AND balance > 0)
OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
""")

Action sur un enregistrement invalide

Vous devez spécifier une action pour déterminer ce qui se passe lorsqu'un enregistrement échoue à la vérification de validation. Le tableau suivant décrit les actions disponibles :

Action

Syntaxe SQL

Syntaxe Python

Résultat

avertir (default)

EXPECT

dp.expect

Les enregistrements non valides sont écrits dans la cible.

déposer

EXPECT ... ON VIOLATION DROP ROW

dp.expect_or_drop

Les enregistrements non valides sont supprimés avant que les données ne soient écrites dans la cible. Le nombre d'enregistrements supprimés est consigné avec d'autres métriques de dataset.

échec

EXPECT ... ON VIOLATION FAIL UPDATE

dp.expect_or_fail

Les enregistrements non valides empêchent la mise à jour d'aboutir. Une intervention manuelle est requise avant le retraitement.

Action

Syntaxe SQL

Syntaxe Python

Résultat

avertir (default)

EXPECT

dp.expect

Les enregistrements non valides sont écrits dans la cible.

déposer

EXPECT ... ON VIOLATION DROP ROW

dp.expect_or_drop

Les enregistrements non valides sont supprimés avant que les données ne soient écrites dans la cible. Le nombre d'enregistrements supprimés est consigné avec d'autres métriques de dataset.

échec

EXPECT ... ON VIOLATION FAIL UPDATE

dp.expect_or_fail

Les enregistrements non valides empêchent la mise à jour d'aboutir. Une intervention manuelle est requise avant le retraitement.

Vous pouvez également implémenter une logique avancée pour mettre en quarantaine les enregistrements non valides sans échouer ni supprimer de données. Voir Mettre en quarantaine les enregistrements non valides.

Métriques de suivi des attentes

Vous pouvez consulter les métriques de suivi des actions warn ou drop à partir de l'interface utilisateur du pipeline. Puisque fail provoque l'échec de la mise à jour lorsqu'un enregistrement non valide est détecté, les métriques ne sont pas enregistrées.

remarque

Pour les tables de streaming et les vues matérialisées soutenues par un pipeline autonome créé dans Databricks SQL, l'onglet Qualité des données dans l'interface utilisateur du pipeline n'est pas disponible. Interrogez les Logs d'événements pour afficher les métriques d'attente. Consultez les métriques de query ou d'attentes de qualité des données.

Pour afficher les métriques d'attente, suivez les étapes suivantes :

  1. Dans la barre latérale de votre workspace Databricks, cliquez sur Tâches & Pipelines .
  2. Cliquez sur le Nom de votre pipeline.
  3. Cliquez sur un dataset avec une attente définie.
  4. Sélectionnez l'onglet Qualité des données dans la barre latérale droite.

Vous pouvez afficher les métriques de qualité des données en interrogeant le Log des événements du LakeFlow Pipelines. Voir Query de la qualité des données ou des métriques d'attentes.

Conserver les enregistrements non valides

La conservation des enregistrements invalides est le comportement par default pour les attentes. Utilisez l'opérateur expect lorsque vous souhaitez conserver les enregistrements qui violent l'attente mais collecter des métriques sur le nombre d'enregistrements qui respectent ou non une contrainte. Les enregistrements qui ne respectent pas l'attente sont ajoutés au dataset cible avec les enregistrements valides :

Python
@dp.expect("valid timestamp", "timestamp > '2012-01-01'")

Supprimer les enregistrements non valides

Utilisez l'opérateur expect_or_drop pour empêcher tout traitement ultérieur des enregistrements non valides. Les enregistrements qui violent l'attente sont supprimés du dataset cible :

Python
@dp.expect_or_drop("valid_current_page", "current_page_id IS NOT NULL AND current_page_title IS NOT NULL")

Échec sur les enregistrements non valides

Lorsque les enregistrements non valides sont inacceptables, utilisez l'opérateur expect_or_fail pour arrêter immédiatement l'exécution lorsqu'un enregistrement échoue à la validation. Si l'opération est une mise à jour de table, le système annule la transaction de manière atomique :

Python
@dp.expect_or_fail("valid_count", "count > 0")
important

Dans un Trigger pipeline, l'échec d'un seul flux n'entraîne pas l'échec des autres flux parallèles. Dans un pipeline continu, l'échec d'une attente arrête le flux, tous les flux dépendants, et le pipeline émet un message expliquant pourquoi il s'est arrêté.

Pour un meilleur contrôle de l’orchestration des workflows lorsqu’une validation échoue, divisez la validation et le travail en aval en pipelines distincts et coordonnez-les avec un flux de contrôle entre les tâches du pipeline. Consultez Tables de validation et flux de contrôle de pipeline.

Graphe d’explication de l’échec de flux LFP

Dépannage des échecs de mise à jour à partir des attentes

Lorsqu'un pipeline échoue en raison d'une violation d'attente, vous devez corriger le code du pipeline pour gérer correctement les données non valides avant de réexécuter le pipeline.

Les attentes configurées pour faire échouer les pipelines modifient le plan de query Spark de vos transformations afin de suivre les informations nécessaires pour détecter et signaler les violations. Vous pouvez utiliser ces informations pour identifier quel enregistrement d'entrée a entraîné la violation pour de nombreuses query. LakeFlow Pipelines fournissent un message d'erreur dédié pour signaler de telles violations. Voici un exemple de message d'erreur de violation d'attente :

Console
[EXPECTATION_VIOLATION.VERBOSITY_ALL] Flow 'sensor-pipeline' failed to meet the expectation. Violated expectations: 'temperature_in_valid_range'. Input data: '{"id":"TEMP_001","temperature":-500,"timestamp_ms":"1710498600"}'. Output record: '{"sensor_id":"TEMP_001","temperature":-500,"change_time":"2024-03-15 10:30:00"}'. Missing input data: false

Gérer plusieurs attentes

remarque

Alors que SQL et Python prennent tous deux en charge plusieurs attentes dans un seul dataset, seul Python vous permet de regrouper plusieurs attentes et de spécifier des actions collectives.

LFP avec Graphe de flux à plusieurs attentes

Vous pouvez regrouper plusieurs attentes et spécifier des actions collectives à l’aide des fonctions expect_all, expect_all_or_drop et expect_all_or_fail.

Ces décorateurs acceptent un dictionnaire Python comme argument, où la clé est le nom de l'attente et la valeur est la contrainte d'attente. Vous pouvez réutiliser le même ensemble d'attentes dans plusieurs dataset de votre pipeline. Voici des exemples de chacun des expect_all opérateurs Python :

Python
valid_pages = {"valid_count": "count > 0", "valid_current_page": "current_page_id IS NOT NULL AND current_page_title IS NOT NULL"}

@dp.table
@dp.expect_all(valid_pages)
def raw_data():
# Create a raw dataset

@dp.table
@dp.expect_all_or_drop(valid_pages)
def prepared_data():
# Create a cleaned and prepared dataset

@dp.table
@dp.expect_all_or_fail(valid_pages)
def customer_facing_data():
# Create cleaned and prepared to share the dataset

Limitations

  • Étant donné que seules les tables en streaming, les vues matérialisées et les vues temporaires prennent en charge les attentes, les métriques de qualité des données ne sont prises en charge que pour ces types d’objets.

  • Les métriques de qualité des données ne sont pas disponibles lorsque :

    • Aucune attente n’est définie sur une query.
    • Un flux utilise un opérateur qui ne prend pas en charge les attentes.
    • Le type de flux, tel que les sinks, ne prend pas en charge les attentes.
    • Il n'y a pas de mises à jour de la table de streaming associée ou de la vue matérialisée pour une exécution de flux donnée.
    • La configuration du pipeline n'inclut pas les paramètres nécessaires pour la capture de métriques, tels que pipelines.metrics.flowTimeReporter.enabled.
  • Dans certains cas, un flux COMPLETED peut ne pas contenir de métriques. Au lieu de cela, les métriques sont signalées dans chaque micro-batch dans un événement flow_progress avec l'état RUNNING.

  • Comme les vues ne sont calculées que lors de la query, les métriques de qualité des données peuvent ne pas être disponibles pour une vue définie. Une autre option serait qu'une vue interrogée dans plusieurs datasets en aval puisse avoir plusieurs ensembles de métriques de qualité des données.

  • Les attentes ne sont pas prises en charge avec AUTO CDC FROM SNAPSHOT.