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.expect— journalise et garde la ligne (surveiller).@dp.expect_or_drop— écarte la mauvaise ligne.@dp.expect_or_fail— interrompt 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 :
- 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 ? ».
- 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).
- Gouvernance centralisée — avec la propagation automatique des permissions
MANAGEaux 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/failobligent 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
- Choisissez une table silver critique et listez 3 à 5 invariants que vous savez devoir tenir.
- 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. - Promouvez les règles stables en
expect_or_drop; réservezexpect_or_failaux quelques invariants qui justifient d'arrêter le chargement. - 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
ai_parse_document() : transformez un PDF en table gouvernée avec une seule instruction SQL
Comment ai_parse_document() de Databricks réunit OCR, parsing et reconstruction de tables en une seule instruction SQL — en livrant le résultat comme table gouvernée dans Unity Catalog.
Lire l'articleChargement incrémental dans Azure Data Factory : le pattern de watermark pas à pas
Comment faire un chargement incrémental dans Azure Data Factory avec le pattern de watermark : Lookup de la dernière valeur, Copy Data uniquement de la nouvelle fenêtre et Stored Procedure qui met à jour le contrôle. Guide pratique.
Lire l'articleVous avez aimé ? Découvrez les e-books pour du contenu approfondi.
E-books