Que sont les Lakeflow Pipelines ?
Les Lakeflow Pipelines fournissent un cadre déclaratif pour la création de pipelines de données batch et de streaming en SQL et Python. Leurs concepts fondamentaux sont les pipelines, les flux, les tables de streaming, les vues matérialisées et les récepteurs, qui fonctionnent ensemble pour traiter les données avec une orchestration automatique et des mises à jour incrémentielles.
LakeFlow Pipelines étendent les Pipelines déclaratifs Apache Spark™ (SDP). Pour en savoir plus sur SDP et sa comparaison avec les LakeFlow pipelines, consultez Apache Spark Declarative Pipelines.
Quels sont les avantages des pipelines ?
Contrairement au développement de processus de data engineering avec les APIs Apache Spark et Spark Structured Streaming sur Databricks Runtime utilisant l’orchestration manuelle via Lakeflow Jobs, la nature déclarative des pipelines offre les avantages suivants :
- **Orchestration automatique** : Les pipelines exécutent les étapes de traitement (appelées « flux ») dans le bon ordre avec un parallélisme maximal, et relancent les échecs transitoires de manière progressive — de la tâche Spark, au flux, à l'ensemble du pipeline.
- Traitement déclaratif : Les fonctions déclaratives réduisent des centaines de lignes de code manuel Spark et Structured Streaming en quelques-unes. L' API AUTO CDC gère les événements de capture de données de changement (CDC)—y compris le SCD Type 1 et Type 2—sans code manuel pour les événements hors séquence ou les concepts de streaming comme les filigranes.
- Traitement incrémentiel : Un moteur de traitement incrémentiel maintient les vues matérialisées à jour : vous rédigez la logique de transformation avec une sémantique par batch, et le moteur retraite uniquement les données source nouvelles ou modifiées lorsque cela est possible.
Concepts clés
Le diagramme ci-dessous illustre les concepts les plus importants des pipelines.

dataset
Un pipeline produit trois types de datasets, chacun avec une sémantique de traitement différente :
Type de dataset | Comment les enregistrements sont traités |
|---|---|
Table de streaming | Chaque enregistrement est traité une seule fois, en supposant une source en mode ajout uniquement. Les tables de streaming sont adaptées à l'ingestion et au traitement incrémentiel des données en croissance continue. |
Vue matérialisée | Les résultats sont recalculés au besoin afin de refléter l'état actuel des données. Les vues matérialisées conviennent aux Transformations, aux agrégations ou aux calculs préalables des résultats consommés par plusieurs datasets en aval. |
Afficher | Évalué à la demande, non persistant. Utilisez des vues pour les transformations intermédiaires et les vérifications qui n'ont pas besoin d'être publiées dans un catalogue. |
Une table de streaming est une forme de table gérée par Unity Catalog qui est également une cible de streaming. Une table de streaming peut contenir un ou plusieurs flux de streaming ( Append , AUTO CDC ) qui y sont écrits. Vous pouvez définir des flux de streaming explicitement et séparément de leur table de streaming cible, ou implicitement dans le cadre d'une définition de table de streaming.
Une *vue matérialisée* est également une forme de table gérée par Unity Catalog et est une cible de traitement par batch. Une vue matérialisée peut contenir un ou plusieurs flux de vues matérialisées. Les vues matérialisées diffèrent des tables de streaming en ce que vous définissez toujours implicitement les flux dans le cadre de la définition de la vue matérialisée.
Pour plus de détails, consultez les tables de streaming et les vues matérialisées.
Quand utiliser les vues, les vues matérialisées et les tables de streaming
Lors de l'implémentation de requêtes de pipeline, choisissez le type de dataset qui correspond le mieux à votre cas d'utilisation.
Envisagez d'utiliser une vue pour :
- Divisez une requête volumineuse ou complexe en requêtes plus faciles à gérer.
- Validez les résultats intermédiaires à l'aide d'attentes.
- Réduisez les coûts de stockage et de compute pour les résultats que vous n'avez pas besoin de conserver. Comme les tables sont matérialisées, elles nécessitent des ressources de compute et de stockage supplémentaires.
Envisagez d'utiliser une vue matérialisée lorsque :
- Plusieurs queries en aval consomment la table. Comme une vue matérialisée met en cache ses résultats, les queries en aval lisent les résultats précalculés au lieu de recalculer la query à chaque accès.
- D'autres pipelines, jobs ou queries consomment la table. Comme une vue matérialisée est matérialisée dans une table Unity Catalog, les consommateurs extérieurs au pipeline qui la définit peuvent la query. Les vues ne sont pas matérialisées, vous ne pouvez donc les utiliser qu’au sein du même pipeline.
- Vous souhaitez inspecter les résultats d'une query pendant le développement. Comme une vue matérialisée est matérialisée et peut être interrogée en dehors du pipeline, vous pouvez valider la justesse des calculs pendant le développement. Après validation, convertissez les requêtes qui ne nécessitent pas de matérialisation en vues.
- Votre query effectue des agrégations ou des jointures, ou les données source peuvent changer en raison de mises à jour et de suppressions plutôt que de simplement croître. Une vue matérialisée maintient ses résultats cohérents avec l'état actuel des données source, tandis qu'une table de streaming est conçue pour les sources à ajout uniquement et traite chaque enregistrement une seule fois.
Envisagez d'utiliser une table de streaming lorsque :
- Une query est définie par rapport à une source de données qui s'accroît continuellement ou de manière incrémentielle.
- Les résultats de la query devraient être compute de manière incrémentale.
- Le pipeline nécessite un throughput élevé et une faible latence.
Les tables de streaming sont toujours définies par rapport aux sources de streaming. Vous pouvez également utiliser des sources de streaming avec AUTO CDC ... INTO pour appliquer les mises à jour des flux CDC. Consultez Les APIs AUTO CDC : Simplifier la capture des données de modification avec des pipelines.
Flux
Un flux est le concept fondamental de traitement des données dans les pipelines et prend en charge les sémantiques streaming et batch. Un flux lit les données d'une source, applique une logique de traitement définie par l'utilisateur et écrit le résultat dans une cible. Les pipelines partagent le même type de flux de streaming ( Append , Update , Complete ) que Spark Structured Streaming. (Actuellement, seuls les flux Ajouter et Mettre à jour sont exposés.) Pour plus de détails, consultez les modes de sortie dans Structured Streaming.
Les pipelines offrent également des types de flux supplémentaires :
- AUTO CDC est un flux de streaming unique dans les LakeFlow pipelines qui gère les événements CDC désordonnés et prend en charge les types SCD 1 et SCD 2. Auto CDC n'est pas disponible dans SDP.
- Une vue matérialisée est un flux de batch dans les pipelines qui ne traite les nouvelles données et les modifications des tables source que lorsque cela est possible.
Pour plus de détails, consultez Charger et traiter les données de manière incrémentielle avec LakeFlow Pipelines.
Puits
Un sink est une cible de streaming pour un pipeline et prend en charge les tables Delta, les topics Apache Kafka, les topics Azure EventHubs et les sources de données Python personnalisées. Un récepteur peut contenir un ou plusieurs flux de streaming (Ajout, Mise à jour) qui y sont écrits.
Pour plus de détails, consultez Dépôts dans les Lakeflow pipelines.
pipeline
Un pipeline est l'unité de développement et d'exécution, et est le conteneur des flux, des tables de streaming, des vues matérialisées et des récepteurs que vous définissez. Vous créez un pipeline en définissant ces objets dans le code source de votre pipeline, puis en exécutant le pipeline. Pendant que votre pipeline s'exécute, il analyse les dépendances de vos objets définis et orchestre automatiquement leur ordre d'exécution et leur parallélisation.
Pour plus de détails, voir Qu’est-ce qu’un pipeline ?.
Vous pouvez également définir des vues matérialisées et des tables de streaming autonomes en dehors d'un Lakeflow pipeline, où Databricks gère le pipeline pour vous. Pour comparer les deux approches, consultez Pipelines autonomes vs. Lakeflow pipelines.
Ingestion des données
Les pipelines prennent en charge toutes les sources de données disponibles dans Databricks. Databricks recommande d'utiliser des tables de streaming pour la plupart des cas d'utilisation d'ingestion. Pour les fichiers stockés dans un stockage d'objets cloud, Auto Loader fournit un chargement incrémentiel et idempotent. Pour les données de streaming, les pipelines peuvent ingérer directement à partir de bus de messages tels que Apache Kafka, Azure Event Hub, Amazon Kinesis et Google Pub/Sub. Voir Charger des données dans les pipelines.
Qualité des données
Les attentes sont des clauses facultatives sur les datasets qui valident les données à mesure qu'elles circulent dans le pipeline. Vous définissez une attente comme une contrainte booléenne SQL et spécifiez ce qui se produit lorsqu'un enregistrement échoue : avertir, ignorer l'enregistrement ou faire échouer la mise à jour. Consultez Gérer la qualité des données avec les attentes de pipeline.
Intégration Delta
Toutes les tables créées et gérées par les pipelines sont des tables Delta. Ils offrent les mêmes garanties que Delta Lake, y compris les transactions ACID, le time travel et l’application des schémas. Les pipelines ajoutent des propriétés de table supplémentaires et effectuent une maintenance automatique à l'aide de l'optimisation prédictive, y compris les opérations OPTIMIZE et VACUUM. Consultez Qu’est-ce que Delta Lake dans Databricks ?.