Aller au contenu principal

Chargez et traitez les données de manière incrémentielle avec LakeFlow Pipelines

Les données sont traitées dans des pipelines par le biais de flux . Chaque flux se compose d'une query et, généralement, d'une cible . Le flux traite la query, soit par batch, soit de manière incrémentielle comme un Stream de données vers la cible. Un flux se trouve dans un Lakeflow pipeline.

Généralement, les flux sont définis automatiquement lorsque vous créez une query dans un pipeline qui met à jour une cible, mais vous pouvez également définir explicitement des flux supplémentaires pour un traitement plus complexe, comme l'ajout à une seule cible à partir de plusieurs sources.

Mises à jour

Un flux est exécuté chaque fois que son pipeline de définition est mis à jour. Le flux va créer ou mettre à jour des tables avec les données les plus récentes disponibles. Selon le type de flux et l'état des modifications apportées aux données, la mise à jour peut effectuer un refresh incrémental, qui ne traite que les nouveaux enregistrements, ou effectuer un refresh complet, qui retraites tous les enregistrements de la source de données.

default flows et flux d'ajout

Lorsque vous créez une query dans un pipeline qui met à jour une cible, un default flow est défini automatiquement. Pour une table en streaming, le flux default est un flux d' ajout qui ajoute de nouvelles lignes à chaque mise à jour, et il porte le même nom que la cible. Créer un flux et sa cible en une seule étape est le moyen le plus courant d'utiliser les pipelines, et vous pouvez l'utiliser pour importer ou transformer des données.

Vous pouvez également définir des flux séparément d'une cible, ce qui permet à plusieurs flux d'ajouter des données à une seule cible. Ceci est utile lorsque vous avez besoin de :

  • Ajoutez des sources de streaming qui s'ajoutent à une table de streaming existante sans nécessiter une refresh complète.
  • Remplir rétroactivement une table de streaming avec les données historiques manquantes.
  • Combiner les données de plusieurs sources sans utiliser une clause UNION.

Pour des exemples de création de flux « default » et explicites, consultez Utiliser les flux dans les LakeFlow Pipelines.

Types de flux

Les flux default pour les tables de streaming et les vues matérialisées sont des flux d’ajout. Vous pouvez également créer des flux pour lire à partir de *sources de données de capture de données modifiées*. Le tableau suivant décrit les différents types de flux.

Type de flux

Description

Ajouter

Les flux d' ajout sont le type de flux le plus courant, où les nouveaux enregistrements dans la source sont écrits vers la cible à chaque mise à jour. Ils correspondent au mode ajout dans le streaming structuré. Vous pouvez ajouter l'indicateur ONCE, indiquant une query batch dont les données doivent être insérées dans la cible une seule fois, sauf si la cible est entièrement actualisée. N'importe quel nombre de flux d'ajout peut écrire sur une cible particulière.

Les flux default (créés avec la table de streaming cible ou la vue matérialisée) auront le même nom que la cible. D'autres cibles n'ont pas de flux default.

Auto CDC (précédemment appliquer les modifications )

Un flux Auto CDC ingère une query contenant des données de capture de changements (CDC). Les flux Auto CDC ne peuvent cibler que les tables de streaming, et la source doit être une source de streaming (même dans le cas des flux ONCE). Plusieurs flux Auto CDC peuvent cibler une seule table de streaming. Une table de streaming qui sert de cible pour un flux Auto CDC ne peut être ciblée que par d'autres flux Auto CDC.

Pour plus d'informations sur les données CDC, consultez Les APIs AUTO CDC : simplifiez la capture des modifications de données avec des pipelines.

Mise à jour (Aperçu public)

Les flux de mise à jour produisent des agrégats de streaming globaux non filigranés vers un sink, n'émettant que les enregistrements qui ont changé dans chaque batch.

Les flux de mise à jour sont uniquement disponibles en Python. Consultez update_flow.

Type de flux

Description

Ajouter

Les flux d' ajout sont le type de flux le plus courant, où les nouveaux enregistrements dans la source sont écrits vers la cible à chaque mise à jour. Ils correspondent au mode ajout dans le streaming structuré. Vous pouvez ajouter l'indicateur ONCE, indiquant une query batch dont les données doivent être insérées dans la cible une seule fois, sauf si la cible est entièrement actualisée. N'importe quel nombre de flux d'ajout peut écrire sur une cible particulière.

Les flux default (créés avec la table de streaming cible ou la vue matérialisée) auront le même nom que la cible. D'autres cibles n'ont pas de flux default.

Auto CDC (précédemment appliquer les modifications )

Un flux Auto CDC ingère une query contenant des données de capture de changements (CDC). Les flux Auto CDC ne peuvent cibler que les tables de streaming, et la source doit être une source de streaming (même dans le cas des flux ONCE). Plusieurs flux Auto CDC peuvent cibler une seule table de streaming. Une table de streaming qui sert de cible pour un flux Auto CDC ne peut être ciblée que par d'autres flux Auto CDC.

Pour plus d'informations sur les données CDC, consultez Les APIs AUTO CDC : simplifiez la capture des modifications de données avec des pipelines.

Mise à jour (Aperçu public)

Les flux de mise à jour produisent des agrégats de streaming globaux non filigranés vers un sink, n'émettant que les enregistrements qui ont changé dans chaque batch.

Les flux de mise à jour sont uniquement disponibles en Python. Consultez update_flow.

Ressources supplémentaires

Pour plus d'informations sur les flux et leur utilisation, consultez les rubriques suivantes :