Brains Up AnalyticsBRAINSUPAnalytics
DatabricksLakeflowData QualityUnity CatalogPySpark

Expectations dans Lakeflow : la qualité des données comme code, gouvernée dans Unity Catalog

Comment déclarer des règles de qualité à côté de la transformation, choisir entre journaliser, écarter ou échouer, et tout gouverner via Unity Catalog avec des règles versionnées et auditables.

Dans la plupart des pipelines dont j'hérite, la « qualité des données » est un script qui tourne après le chargement : une batterie de SELECT COUNT(*) WHERE champ IS NULL que quelqu'un ouvre le lundi et, quand il trouve un problème, la mauvaise donnée a déjà été consommée par trois tableaux de bord et un modèle. C'est de la réaction, pas de la prévention.

Les Expectations de Lakeflow (les Databricks Declarative Pipelines, évolution de l'ancien DLT) inversent cette logique. La règle de qualité vit désormais à côté de la transformation, comme du code versionné, et Unity Catalog gouverne, audite et réutilise cette règle entre les pipelines. Dans cet article, je montre le concept, le code en pratique et pourquoi c'est important pour qui exploite des données en production.

Le concept : la qualité comme contrat, pas comme vérification

Une Expectation est une contrainte déclarative sur un dataset de pipeline — une materialized view, une streaming table ou une vue temporaire. Vous décrivez, en SQL booléen, ce que signifie un enregistrement valide et vous dites au pipeline quoi faire quand la condition n'est pas satisfaite.

Le changement de mentalité tient dans le mot contrat. Au lieu de « je vérifierai plus tard si c'est arrivé nul », vous déclarez « cette table n'accepte pas d'identifiant nul » et le moteur s'occupe du reste — y compris d'enregistrer combien d'enregistrements ont violé la règle, pour que vous mesuriez la santé de la source dans le temps.

Le code en pratique

Le point d'entrée en 2026 est le module pipelines de PySpark. Vous l'importez, décorez la fonction qui produit la table et attachez l'expectation :

from pyspark import pipelines as dp

@dp.table()
@dp.expect_or_drop("id_valide", "id IS NOT NULL")
def clients():
    return spark.readStream.table("bronze.clients_raw")

Trois éléments ici :

  • @dp.table() déclare que la fonction matérialise une table du pipeline.
  • @dp.expect_or_drop(...) attache la règle. Le premier argument est le nom de l'expectation (il apparaît dans les métriques), le second est la condition SQL qu'un enregistrement valide doit satisfaire.
  • La fonction renvoie la requête — Spark gère le plan d'exécution.

Les 6 décorateurs et quand utiliser chacun

Le module offre six décorateurs, organisés par ce qui arrive à la violation et par nombre de règles.

Règle unique :

  • @dp.expectjournalise et garde la ligne (surveiller).
  • @dp.expect_or_dropécarte la mauvaise ligne.
  • @dp.expect_or_failinterrompt le pipeline.

Plusieurs règles (dictionnaire) :

  • @dp.expect_all — journalise et garde.
  • @dp.expect_all_or_drop — écarte.
  • @dp.expect_all_or_fail — interrompt.

Le choix est une décision métier, pas une décision de code :

  • expect — à utiliser quand vous voulez de l'observabilité sans bloquer. Idéal pour commencer : vous mesurez le taux de violation d'une nouvelle source avant de décider de bloquer.
  • expect_or_drop — à utiliser quand la mauvaise ligne ne peut pas atteindre la consommation, mais que le pipeline peut continuer avec le reste. C'est le cas le plus courant dans les couches silver/gold.
  • expect_or_fail — à utiliser pour les invariants critiques du métier (ex. : clé primaire dupliquée, valeur monétaire négative là où c'est impossible). Échouer tôt coûte moins cher que propager.

Plusieurs règles d'un coup

Pour valider plusieurs conditions, utilisez les variantes _all, en passant un dictionnaire nom -> condition :

regles = {
    "id_valide":       "id IS NOT NULL",
    "age_plausible":   "age BETWEEN 0 AND 120",
    "region_remplie":  "region IS NOT NULL",
}

@dp.table()
@dp.expect_all_or_drop(regles)
def clients():
    return spark.readStream.table("bronze.clients_raw")

Avec expect_all_or_drop, un enregistrement est écarté s'il échoue à n'importe laquelle des règles. Chaque règle reste mesurée individuellement dans les métriques — vous savez quelle condition fait tomber les données.

Le saut de 2026 : des règles gouvernées par Unity Catalog

Historiquement, les expectations vivaient attachées au code du pipeline. La nouveauté qui change la donne, c'est de pouvoir stocker et gérer les règles de qualité dans des tables d'Unity Catalog. En pratique, cela apporte trois bénéfices directs :

  1. Versionnage et audit — la règle cesse d'être une ligne perdue dans un notebook et devient un objet gouverné, avec un historique. L'audit et la conformité peuvent répondre « quelle règle était en vigueur en mars ? ».
  2. Réutilisation entre pipelines — la même définition de « client valide » peut être référencée par plusieurs pipelines, au lieu d'être copiée-collée (et de diverger avec le temps).
  3. Gouvernance centralisée — avec la propagation automatique des permissions MANAGE aux materialized views et streaming tables dans Unity Catalog, la qualité entre dans le même périmètre de gouvernance que le reste des données.

Ajoutez à cela d'autres évolutions récentes de Lakeflow qui rendent le pattern plus robuste en production : des tests unitaires en Python directement dans l'éditeur de pipelines (validant la logique contre des données mockées par redirection de table), le type widening pour faire évoluer les types de colonne sans reset du pipeline, et un mode d'exécution en file qui met en file les updates concurrents au lieu d'échouer sur conflit.

Pourquoi c'est important

Quand la qualité devient du code déclaratif et gouverné, trois choses changent pour l'équipe :

  • Des données fiables sans surveiller le chargement. La règle tourne à chaque exécution ; la violation est journalisée automatiquement. Personne n'a besoin d'ouvrir le tableau de bord le lundi matin.
  • Une décision de risque explicite. expect / drop / fail obligent l'équipe à décider, pour chaque règle, le coût de laisser passer versus le coût de bloquer. Cela documente l'appétit pour le risque dans le pipeline lui-même.
  • Une observabilité native. Les métriques de qualité font partie du pipeline — on peut suivre le taux de violation par source dans le temps et agir avant que ça ne devienne un incident.

Comment commencer demain

  1. Choisissez une table silver critique et listez 3 à 5 invariants que vous savez devoir tenir.
  2. Commencez avec @dp.expect (surveiller seulement) pendant quelques jours pour mesurer le taux réel de violation — cela évite d'écarter de la bonne donnée à cause d'une règle mal calibrée.
  3. Promouvez les règles stables en expect_or_drop ; réservez expect_or_fail aux quelques invariants qui justifient d'arrêter le chargement.
  4. Quand le pattern mûrit, déplacez les définitions dans Unity Catalog et partagez-les entre pipelines.

La qualité cesse d'être une réaction hebdomadaire et devient un contrat que le pipeline honore à chaque exécution. En pratique, c'est le type de changement qui paie des dividendes plus vite que presque n'importe quelle optimisation de performance — parce qu'il évite le retravail silencieux de retraiter une mauvaise donnée déjà consommée.

Articles liés

Vous avez aimé ? Découvrez les e-books pour du contenu approfondi.

E-books