Comment utiliser les Lakeflow pipelines
Cette page explique comment utiliser les Lakeflow pipelines tout au long du cycle de vie d’un pipeline de données, des premières décisions de conception jusqu’à l’exécution à grande échelle, ainsi que les compromis associés à chaque étape. Chaque section contient des Link vers les articles qui vous expliquent comment faire.
Ce guide suppose une familiarité avec les concepts fondamentaux du data engineering. Si vous débutez avec les pipelines, commencez par Apache Spark Declarative Pipelines pour découvrir ce qu'est le produit et le modèle déclaratif sous-jacent, puis suivez le Tutoriel : Construire un pipeline ETL à l'aide de la capture de données modifiées (CDC).
Vue d'ensemble du cycle de vie du pipeline
Un pipeline se déroule en six étapes :
- Planifiez et concevez: Décidez ce que vous construisez et choisissez les outils, le langage et compute qui vous conviennent.
- Ingérer des données: importez les données sources dans le pipeline de manière fiable et incrémentielle.
- Transformer et modéliser: nettoyer, valider, relier et façonner les données en tables fiables pour les consommateurs.
- Opérationnaliser: placez le pipeline sous contrôle de version, testez-le, planifiez-le et promouvez-le dans les différents environnements.
- Exécuter en production: surveillez, alertez, déboguez, effectuez des remplissages, sécurisez et suivez la traçabilité pendant que le pipeline s'exécute sans surveillance.
- Maturer et monter en charge: confirmer la préparation à la production et maintenir la santé du pipeline à mesure que le volume et la taille de l'équipe augmentent.
Les étapes ne sont pas strictement séquentielles, mais elles correspondent à l'ordre dans lequel les questions se posent. Comme les Lakeflow pipelines gèrent l'orchestration, les points de contrôle, les nouvelles tentatives et le traitement incrémentiel, votre travail à chaque étape relève davantage d'une décision de conception que de la mise en œuvre.
Planification et conception
Vos premières décisions façonnent tout en aval. Pour savoir comment le modèle déclaratif se compare à l’écriture d’étapes procédurales vous-même, voir Traitement procédural vs. déclaratif de données dans Databricks.
Quelques choix définissent votre configuration de départ :
- Un dataset autonome ou un pipeline. Une vue matérialisée ou une table de streaming unique peut être définie en SQL comme un dataset autonome, et Databricks gère le pipeline de refresh en arrière-plan. Créez et exploitez un pipeline Lakeflow en tant qu'unité lorsque vous avez besoin de la création en Python, de récepteurs (sinks) ou d'une orchestration multi-étapes. Voir Pipelines autonomes vs. LakeFlow Pipelines.
- SQL ou Python (ou les deux). SQL convient aux transformations qui sont principalement des filtres, des jointures et des agrégations. Python convient à la logique personnalisée, aux bibliothèques externes ou à la génération de nombreuses tables similaires de manière programmatique. Le choix se fait par fichier plutôt que sur tout le pipeline, donc vous pouvez mélanger les deux sans avoir à régler tout d’avance.
- Compute serverless ou classic. Le mode Serverless est le choix par default recommandé et supprime la configuration de cluster. Choisissez le mode classic lorsque vous avez besoin de types d'instances spécifiques, de politiques de cluster personnalisées ou d'un script d'initialisation. Voir Configurer un pipeline serverless et Configurer un compute classic pour les pipelines.
- Exécution déclenchée ou continue. Start Trigger, since it only consumes compute while it runs. Le mode continu permet de faire tourner le compute pour traiter les nouvelles données avec un délai minimal, ce qui est généralement le facteur de coût le plus important, donc réservez-le à une exigence de latence prouvée. Voir mode pipeline déclenché vs. mode pipeline continu.
Un pipeline déduit son Graphe d'exécution à partir des datasets référencés par votre code ; le travail de conception consiste donc principalement à nommer et à séquencer les datasets. La décision principale consiste à choisir le type de chaque sortie : une table de streaming pour les données incrémentielles à ajout intensif, ou une vue matérialisée pour les agrégats et les jointures recalculés. Ce choix influe sur le coût et l'exactitude, car le traitement incrémentiel monte en charge avec le taux de nouvelles données, tandis qu'un recalcul complet monte en charge avec l'ensemble de votre historique. Pour savoir quel type correspond à quel Job, consultez Qu'est-ce qu'un pipeline ?.
Comme le code pipeline est du Python et du SQL ordinaires, vous pouvez l’écrire, l’utiliser et le valider dans votre propre éditeur avant de le déployer dans un Workspace partagé.
À cette étape
Questions à considérer pendant cette étape :
- Comment choisir entre un dataset autonome et un pipeline complet ?
- Comment puis-je identifier mes sources de données et comprendre comment m’y connecter ?
- Comment puis-je concevoir l’architecture de mon pipeline avant d’écrire du code ?
- Comment choisir un format de fichier et une couche de stockage ?
- Comment configurer un environnement de développement local ?
- Comment planifier la montée en charge et estimer les coûts avant de start à développer ?
Ingérer des données
La question centrale de conception est de savoir si une source est en ajout seul ou si elle est modifiée sur place. Cela détermine la manière dont vous modélisez la cible :
- Sources en ajout seul , telles que les fichiers arrivant dans le stockage cloud ou les événements sur un bus de messages, sont ingérées dans une table de streaming, qui enregistre sa progression afin qu'un redémarrage ne retraite ni ne perde de données. Auto Loader gère les fichiers, en découvrant les nouveaux, tout en déduisant et en faisant évoluer le schéma à mesure qu'ils arrivent. Les bus de messages tels qu'Apache Kafka, Azure Event Hubs, Amazon Kinesis et Google Pub/Sub effectuent une lecture directe dans une table de streaming. Dédupliquez en aval, car un bus peut transmettre le même événement plus d'une fois. Pour Azure Event Hubs spécifiquement, consultez Utiliser Azure Event Hubs comme source de données de pipeline.
- Les sources qui mettent à jour et suppriment des lignes , comme la plupart des bases de données et de nombreux systèmes SaaS (outil/solution/technologie/plateforme as a Service), utilisent la capture des modifications de données (CDC). Une copie complète à chaque exécution est inefficace et devient plus lente à mesure que la source augmente ; le CDC ne lit donc que les lignes ayant changé depuis la dernière exécution. L'API
AUTO CDCapplique ces modifications sans logique de Merge manuelle ; voir Les APIs AUTO CDC : simplifiez la capture des modifications de données avec les pipelines. Un flux applique le CDC dans une table de streaming, et plusieurs flux peuvent alimenter une seule table, ce qui vous permet de regrouper plusieurs sources vers une cible unique.
La création de points de contrôle et les réessais sont automatiques, de sorte qu'un pipeline reprend à partir du dernier offset traité plutôt que de tout retraiter. Deux garde-fous sont activables sur demande :
- Une colonne de données sauvées capture les enregistrements qui ne correspondent pas au schéma attendu.
- Les attentes appliquent l’action au niveau de la ligne que vous définissez.
Si un point de contrôle de streaming devient non valide, privilégiez la récupération la moins coûteuse qui préserve les données de la table.
À cette étape
Questions à considérer pendant cette étape :
- Comment ingérer des données à partir d'une base de données et choisir entre un chargement complet et la CDC ?
- Comment ingérer des données à partir d’une API ?
- Comment ingérer des données de streaming ou des données d'événement ?
- Comment ingérer des fichiers de manière fiable ?
- Comment gérer les échecs d’ingestion sans perdre de données ?
Transformer et modéliser
La transformation convertit les données ingérées en tables propres, fiables pour les utilisateurs et les outils. C’est ici que le modèle médaillon (bronze à argent à or) prend une forme concrète.
Le nettoyage et la validation passent avant tout. Les attentes sont une fonctionnalité intégrée du pipeline Lakeflow : des contraintes de qualité des données que le pipeline évalue sur chaque ligne de chaque exécution, en indiquant le nombre de réussites et d’échecs, de sorte que la qualité est continue plutôt qu’un contrôle ponctuel. Décidez ce qui se passe lorsqu’une ligne échoue (prévenez-la et gardez-la, supprimez-la ou échouez la mise à jour) et où la porte doit être. Les portes se situent généralement à la limite bronze-argent, donc tout ce qui est en aval peut être fiable sans vérifier à nouveau.
La jointure et l’agrégation façonnent l’étape argent-vers-or. Une vue matérialisée convient à une jointure ou une agrégation de type batch sur des tables existantes, car elle maintient les résultats cohérents avec ses sources : elle refresh de manière incrémentielle lorsque la query et les sources le permettent, et sinon, elle recalcule en totalité, produisant le même résultat dans les deux cas. Cela en fait le bon choix lorsque l’exactitude importe plus que la latence, car il recalcule les jointures lorsqu’une dimension change. Voir Comment les pipelines refresh ?. La jointure de flux en direct augmente l’état non délimité ; par conséquent, les jointures et agrégations en streaming nécessitent un filigrane pour limiter la durée pendant laquelle le pipeline attend les données arrivant en retard.
Deux idées de correction traversent cette étape :
- L’idempotence signifie qu’un pipeline produit le même résultat, peu importe le nombre de fois où il s’exécute sur la même entrée. LakeFlow Pipelines sont idempotents pour les éléments qu’ils gèrent, tels que les lectures par checkpoint et les surserts de
AUTO CDCbasés sur des clés ; Vous conservez votre propre logique idempotente en évitant les fonctions non déterministes dans les vues recalculées. - Traitement « at-least-once » (au moins une fois) versus « exactly-once » (exactement une fois). Les tables Delta-to-Delta gérées valident (commit) les entrées et sorties de chaque micro-batch ensemble, vous offrant ainsi un traitement « exactly-once » par default. Cela s'arrête aux limites, comme un puits personnalisé, une cible non Delta ou une source personnalisée non vérifiée, où vous traitez l'écriture comme « at-least-once » et la rendez idempotente, par exemple en effectuant un upsert sur une clé.
Les dimensions à évolution lente (SCD) se trouvent également ici : AUTO CDC implémente directement le SCD de type 1 et de type 2, vous définissez donc un type plutôt que d'écrire une logique de suivi de l'historique.
À cette étape
Questions à considérer pendant cette étape :
- Comment nettoyer et valider les données entrantes ?
- Comment suivre l’historique dans le temps avec des dimensions qui changent lentement (SCD) ? Qu’est-ce que la SCD ?
- Comment joindre des données de streaming et des données statiques ? Comment agréger efficacement des données ?
- Comment modéliser mes données pour une utilisation en aval ?
- Comment garantir les traitements dans les Lakeflow pipelines ?
- Traitement « au moins une fois » (at-least-once) versus « exactement une fois » (exactly-once) : quelle est la différence et de quoi ai-je besoin ?
- Comment gérer les données arrivant en retard ou dans le désordre ?
Mettre en service
L’opérationnalisation fait passer un pipeline de quelque chose qui fonctionne pour vous à quelque chose que l’équipe peut construire, tester et livrer de manière répétée. Un pipeline est le code source plus la configuration, donc les pratiques classiques d’ingénierie logicielle s’appliquent.
Les tests couvrent deux aspects simultanément : votre logique de transformation et la qualité continue des données qui y transitent. Les attentes (Expectations) gèrent le volet données en continu. Pour la logique, factorisez les transformations en fonctions simples et testez-les unitairement en dehors du runtime, puis validez le graphe du pipeline avec une simulation (dry run) avant toute matérialisation. Consultez Tests unitaires pour les pipelines.
Gardez le code pipeline dans Git et emballez-le pour le déploiement afin qu’il puisse être examiné, annulé et déployé de manière cohérente à travers les environnements. Le package n’est pas une alternative aux Lakeflow pipelines. C’est le projet et le CI/CD enveloppant votre pipeline, et votre logique de données reste déclarative. Paramétrez les valeurs spécifiques à l’environnement comme les noms de catalogue et les chemins afin que le même code s’exécute sans modification dans chaque environnement. Voir Utiliser des paramètres avec les pipelines.
Pour exécuter un pipeline selon un calendrier, encapsulez-le dans un Run pipelines in a workflow: Databricks recommande de planifier et d'orchestrer les pipelines avec des jobs, ce qui vous permet également de coordonner le pipeline avec d'autres tâches, comme l'enchaînement d'un rapport en aval ou de plusieurs pipelines. Au sein d'une exécution, un pipeline ordonne et parallélise ses propres datasets, de sorte que l'orchestration ne coordonne que les tâches situées en dehors du pipeline.
À cette étape
Questions à considérer pendant cette étape :
- Comment tester un pipeline de données et pourquoi est-ce différent du test d'un logiciel classique ?
- Comment gérer le contrôle de version et collaborer sur le code de pipeline en équipe ?
- Comment puis-je programmer ou orchestrer mon pipeline pour qu’il s’exécute automatiquement ?
- Comment puis-je faire passer mon pipeline du développement à la staging puis à la production en toute sécurité ?
- Comment configurer la CI/CD pour mon pipeline ?
Exécuter en production
Une fois qu'un pipeline s'exécute sans surveillance sur des données réelles, le travail consiste à savoir s'il est sain et à le corriger lorsqu'il ne l'est pas.
Le monitoring fonctionne à trois niveaux de profondeur. La liste Jobs & Pipelines donne un aperçu rapide des dernières expéditions. L’interface de monitoring du pipeline affiche chaque table et flux codés par couleur selon le statut, avec le nombre de lignes, les métriques de qualité des données et les métriques de backlog pour les tables de streaming. L’event log sous les deux est la source de vérité pour tout ce qui est programmatique ou historique. Configurez les notifications de défaillance pour être informé d’une exécution défaillante avant que vos parties prenantes ne la signalent. Pour un aperçu des surfaces de monitoring, voir Surveiller les pipelines.
Déboguez en remontant de l'échec mis en évidence sur le Graphe jusqu'aux détails complets de l'erreur dans le log des événements, puis relancez uniquement ce qui a échoué. Le comportement de réessai diffère selon le trigger : les mises à jour déclenchées manuellement désactivent les réessais automatiques afin que vous voyiez les erreurs immédiatement, tandis que les mises à jour planifiées réessaient en cas d'échecs récupérables. Une alerte de production peut donc se résoudre d'elle-même lors d'une nouvelle tentative, alors que la même erreur ne le fera pas lors d'un développement interactif. Pendant le développement, Genie Code peut vous aider à diagnostiquer et à corriger les erreurs au niveau du code au fur et à mesure de vos itérations, bien qu'il cible aujourd'hui la création de pipelines plutôt que le diagnostic des exécutions en production.
Modélisez un backfill comme son propre flux explicite et ponctuel alimentant la même cible que votre flux incrémentiel habituel. Le garder séparé enregistre quand et comment l'historique a été chargé et maintient la logique d'état stable simple.
Sécuriser un pipeline en contrôlant qui peut l’exploiter, en le faisant fonctionner comme un Service Principal dédié plutôt qu’un compte personnel, et en gardant les identifiants dans un Secret Scope plutôt que dans le code source. La lignée est automatique, capturée jusqu’au niveau de la colonne. Un pipeline écrit vers un système externe via des Sinks dans LakeFlow Pipelines, la limite où la réflexion au moins une fois ci-dessus s’applique.
À cette étape
Questions à considérer pendant cette étape :
- Comment puis-je vérifier si mon pipeline a été exécuté avec succès ?
- Comment être alerté en cas de problème ?
- Comment déboguer une exécution de pipeline ayant échoué ?
- Comment effectuer un remplissage rétroactif des données historiques ?
- Comment puis-je contrôler et prédire le coût d'exécution de mon pipeline ?
- Comment sécuriser mon pipeline, notamment en ce qui concerne les identifiants, le contrôle d'accès et les données PII ?
- Comment documenter mon pipeline et suivre la data lineage ?
Maturité et Monter en charge
Un pipeline mature s'exécute sans surveillance et évolue sans réécriture. La confirmation de l'état de préparation et la planification de la montée en charge définissent cette étape.
La préparation à la production est une liste de contrôle couvrant la qualité des données, la fiabilité, l'observabilité, le déploiement, les coûts et la gouvernance. Considérez chaque élément non coché comme une lacune connue : chaque dataset susceptible de recevoir des données erronées dispose-t-il d'une attente, le pipeline est-il planifié plutôt que démarré manuellement, les notifications d'échec sont-elles configurées, s'exécute-t-il en tant que Service Principal, est-il déployé à partir du contrôle de version sur au moins une cible de développement et de production. La qualité des données et les notifications sont les éléments les moins coûteux à ajouter et les plus susceptibles de détecter une exécution erronée non identifiée.
Monter en charge en réponse à des signaux concrets indiquant que la santé du pipeline se dégrade :
- La durée de mise à jour tend à augmenter.
- Le dimensionnement automatique atteint régulièrement son plafond.
- Le coût augmente plus rapidement que l’entreprise sous-jacente.
- Les vues matérialisées reviennent à des recalculs complets.
Essayez d'abord les leviers au niveau du compute, comme le passage au serverless ou l'adaptation de son mode de performance à vos besoins en matière de latence. Au-delà de cela, la manière dont vous organisez les ensembles de données entre pipelines est la plus importante :
- Un pipeline a une limite de simultanéité : il ne met à jour qu'un nombre défini de datasets en même temps. Une fois qu'un pipeline contient plus de datasets que cette limite, les mises à jour supplémentaires attendent dans une file d'attente, ce qui augmente le temps de mise à jour total du pipeline.
- Regroupez les datasets associés et fractionnez ceux qui ne le sont pas. Regroupez par domaine, cadence de refresh partagée et dépendance ; fractionnez au niveau des limites de propriété, de couche et de latence. La séparation de l'ingestion et de la transformation, par exemple, évite qu'une ingestion lente ne retarde tout le processus en aval et permet de maintenir chaque pipeline suffisamment petit pour rester sous la limite de concurrence.
Il est plus facile de fusionner deux petits pipelines ultérieurement que de diviser un grand pipeline déjà en production. Pour savoir comment regrouper et diviser les datasets, consultez Organiser les datasets dans les LakeFlow Pipelines.
À cette étape
Questions à considérer pendant cette étape :