Aller au contenu

Data Flows — pipelines visuels

Le Data Flow designer assemble un graphe de nœuds : des Sources, des Transforms et des Écritures (Sinks). SAAM le compile en un job Sparkbatch (lecture bornée) ou streaming (Structured Streaming).

Le catalogue de nœuds

Famille Exemples Rôle
Source jdbc · iceberg · kafka · fichier Lire depuis une base, une table du lakehouse, un topic ou un fichier objet.
Transform SQL · filtre · jointure · inférence Transformer, enrichir, joindre — ou appeler un modèle exposé (scoring).
Écriture iceberg · minio · kafka · jdbc Matérialiser le résultat dans une table, un bucket, un topic ou une base.

Procédure

1. Nouveau flux

Nommez-le, choisissez le mode (Batch / Streaming) et la visibilité.

2. Ajouter les nœuds

Boutons Source, Transform, Écriture. Reliez-les en tirant des ports d'un nœud à l'autre. Exemple : une Source JDBC PostgreSQL → un Transform SQL → une Écriture Iceberg (couche bronze).

3. Dimensionner & générer

Réglez les ressources Spark (exécuteurs, mémoire), prévisualisez le Code Spark généré (‹/› Code Spark), puis Déployer ().

4. Exécuter & suivre

Démarrer () lance le job ; l'onglet Runs () montre l'historique et l'état d'exécution.

flowchart LR
  S1[(Transactions PG)]:::src --> T[SQL / jointure]:::tf
  S2[(Clients PG)]:::src --> T
  T --> W[(Bronze Iceberg)]:::sink
  classDef src fill:#eef6f2,stroke:#17987a
  classDef tf fill:#fff6e6,stroke:#e3a63c
  classDef sink fill:#eaf0f4,stroke:#0e2233

Batch ou streaming

Lecture bornée (fichier, JDBC, Iceberg, Kafka borné). Idéal pour la collecte historique et les recalculs. Pas de watermark.

Structured Streaming sur des sources continues (Kafka). Gestion du watermark pour l'agrégation temporelle et les jointures.

Cycle de vie d'un flux

  1. Prévisualiser le Code Spark — inspectez à tout moment le job généré à partir du graphe (utile pour comprendre ou auditer la logique).
  2. Versionner — chaque enregistrement crée une version ; comparez et revenez en arrière depuis l'onglet Versions.
  3. Déployer & démarrerDéployer publie le job ; Démarrer le lance (ou le maintient actif en streaming).
  4. Suivre les Runs — historique d'exécutions, état et logs.

Lineage automatique

Un Data Flow enregistre son lineage à chaque déploiement : on sait quelles sources alimentent quelles cibles. C'est la base de la traçabilité côté Gouvernance.

Import / export

Un flux s'exporte en fichier .flow.json et se réimporte sur un autre cluster — pratique pour versionner et rejouer un pipeline.

Suivant : Jobs programmés (Airflow) →