Pipelines de production avec des workflows

Transformation des données avec Spark SQL dans Databricks

Disha Mukherjee

Lead Data Engineer

Pourquoi Delta Lake ?

Un coffre-fort numérique sécurisé et lumineux avec des tables de données organisées et un bouclier de protection, style moderne et épuré

 

$$

  • Transactions ACID → annuler les écritures échouées
  • Respect du schéma → bloquer les types non conformes
  • Gestion des versions → interroger tout état antérieur
Transformation des données avec Spark SQL dans Databricks

Écrire vers Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • Une nouvelle table Delta apparaît dans Unity Catalog
Transformation des données avec Spark SQL dans Databricks

Tâches du carnet

task1_ingest : charger et nettoyer

  • Charge le CSV → applique le nettoyage → écrit une table Delta
Transformation des données avec Spark SQL dans Databricks

Tâches du carnet

task2_metrics : revenus par catégorie

  • Lit la table nettoyée → calcule des mesures → écrit dans une nouvelle table
Transformation des données avec Spark SQL dans Databricks

Tâches du carnet

task3_customers : classement par dépenses

  • Lit la table nettoyée → classe la clientèle → enregistre dans une nouvelle table
Transformation des données avec Spark SQL dans Databricks

Création du travail

Interface Jobs and Pipelines montrant un DAG de tâche à trois étapes terminé avec des flèches de dépendance de task1_ingest à task2_metrics puis à task3_customers

  • Peut s'exécuter manuellement, planifier des tâches et définir des déclencheurs
Transformation des données avec Spark SQL dans Databricks

Exécuter le travail

DAG d'exécution : task1_ingest et task2_metrics en vert (réussite) et task3_customers en rouge (échec)

Transformation des données avec Spark SQL dans Databricks

Exécuter le travail

Vérification d'une erreur dans la troisième tâche

  • Erreur : customer-id devrait être Customer_ID
Transformation des données avec Spark SQL dans Databricks

Exécuter le travail

Vues de graphe avec réussite

Une vue chronologique

Transformation des données avec Spark SQL dans Databricks

Qu'est-ce que Lakeflow ?

$$

$$

comparaison : Impératif (Jobs) | Déclaratif (Lakeflow)

 

  • Jobs → vous gérez chaque étape
  • Lakeflow → déclare ce que les tables doivent contenir
  • Databricks gère l'ordre, les nouvelles tentatives et le calcul
Transformation des données avec Spark SQL dans Databricks

Le modèle @dlt.table

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
Transformation des données avec Spark SQL dans Databricks

Exécution du pipeline

Graphe de pipeline Lakeflow montrant transactions_bronze 100K lignes puis transactions_silver 33K lignes puis category_revenue_gold 6 lignes, toutes des vues matérialisées en vert

  • Bronze (100K lignes) → Silver (33K lignes) → Gold (6 lignes)
Transformation des données avec Spark SQL dans Databricks

Carnets, Jobs ou Lakeflow ?

$$

couches : Carnets, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Carnets → exploration, prototypage
  • Databricks Jobs → pipelines à étapes multiples, planifiés
  • Lakeflow → entièrement géré, déclaratif
Transformation des données avec Spark SQL dans Databricks

Passons à la pratique !

Transformation des données avec Spark SQL dans Databricks

Preparing Video For Download...