Exercice Lakehouse
Contexte
Vous êtes data engineer dans une entreprise de vente en ligne. Vous disposez de fichiers CSV contenant les commandes clients des trois dernières années (2022, 2023, 2024), déposés dans le dossier Files/commandes/ de votre lakehouse Fabric. Votre mission : charger ces données, les nettoyer, les transformer et les stocker sous forme de table Delta interrogeable en SQL.
Prérequis
- Un espace de travail Microsoft Fabric avec une capacité active (essai ou Premium).
- Un lakehouse créé dans cet espace de travail.
- Un notebook Fabric attaché au lakehouse.
Données d’exemple
Chaque fichier CSV (2022.csv, 2023.csv, 2024.csv) contient les colonnes suivantes :
NumeroCommande,LigneCommande,DateCommande,NomClient,Email,Produit,Quantite,PrixUnitaire,Taxe
SO001,1,2022-01-15,Marie Dupont,marie@example.com,Clavier,2,45.99,9.20
SO001,2,2022-01-15,Marie Dupont,marie@example.com,Souris,1,25.50,5.10
SO002,1,2022-01-18,Jean Martin,jean@example.com,Ecran,1,299.00,59.80Étape 1 : charger un fichier CSV dans un dataframe
La commande de base pour lire un fichier CSV dans un dataframe Spark :
# Charger un seul fichier CSV
df = spark.read.format("csv").option("header", True).load("Files/commandes/2022.csv")
# Afficher les premières lignes
display(df)Le paramètre header=True indique que la première ligne du fichier contient les noms de colonnes.
Etape 2 : charger plusieurs fichiers CSV avec un schéma explicite
Pour charger tous les fichiers du dossier et forcer les types de données corrects :
from pyspark.sql.types import *
# Définir le schéma explicitement
schema_commandes = StructType([
StructField("NumeroCommande", StringType()),
StructField("LigneCommande", IntegerType()),
StructField("DateCommande", DateType()),
StructField("NomClient", StringType()),
StructField("Email", StringType()),
StructField("Produit", StringType()),
StructField("Quantite", IntegerType()),
StructField("PrixUnitaire", FloatType()),
StructField("Taxe", FloatType())
])
# Charger tous les CSV du dossier avec le wildcard *
df = spark.read.format("csv") \
.schema(schema_commandes) \
.option("header", True) \
.load("Files/commandes/*.csv")
# Vérifier le nombre de lignes chargées
print(f"Nombre de lignes : {df.count()}")
# Vérifier le schéma
df.printSchema()
# Afficher un échantillon
display(df.limit(10))Etape 3 : explorer les données
Quelques commandes utiles pour explorer rapidement les données chargées :
# Nombre de lignes et de colonnes
print(f"Lignes : {df.count()}, Colonnes : {len(df.columns)}")
# Statistiques descriptives sur les colonnes numériques
display(df.describe())
# Compter les valeurs distinctes d'une colonne
from pyspark.sql.functions import countDistinct
df.select(countDistinct("NomClient").alias("NbClients")).show()
# Vérifier les valeurs nulles par colonne
from pyspark.sql.functions import col, sum as spark_sum
display(
df.select([spark_sum(col(c).isNull().cast("int")).alias(c) for c in df.columns])
)Etape 4 : nettoyer et transformer les données
4.1. Supprimer les doublons
# Supprimer les lignes en double
df_clean = df.dropDuplicates()
print(f"Avant : {df.count()} lignes, Après : {df_clean.count()} lignes")4.2. Supprimer les lignes avec des valeurs nulles sur les colonnes clés
# Supprimer les lignes où NumeroCommande ou DateCommande est null
df_clean = df_clean.dropna(subset=["NumeroCommande", "DateCommande"])4.3. Ajouter des colonnes calculées
from pyspark.sql.functions import col, year, month, quarter, round
# Ajouter le montant total de la ligne
df_clean = df_clean.withColumn(
"MontantLigne",
round((col("PrixUnitaire") * col("Quantite")) + col("Taxe"), 2)
)
# Extraire l'année, le trimestre et le mois de la date de commande
df_clean = df_clean.withColumn("Annee", year(col("DateCommande")))
df_clean = df_clean.withColumn("Trimestre", quarter(col("DateCommande")))
df_clean = df_clean.withColumn("Mois", month(col("DateCommande")))
display(df_clean.limit(5))4.4. Renommer une colonne
# Renommer une colonne
df_clean = df_clean.withColumnRenamed("PrixUnitaire", "PrixUnit")4.5. Filtrer les données
# Ne garder que les commandes de 2023 et après
df_recent = df_clean.filter(col("Annee") >= 2023)
print(f"Commandes depuis 2023 : {df_recent.count()} lignes")Etape 5 : sauvegarder en table Delta dans le lakehouse
5.1. Créer une table Delta (écriture complète)
# Sauvegarder le dataframe comme table Delta dans le lakehouse
df_clean.write.mode("overwrite").format("delta").saveAsTable("commandes")- mode("overwrite") : remplace la table si elle existe déjà.
- format("delta") : utilise le format Delta Lake (versionnement, transactions ACID).
- saveAsTable("commandes") : crée la table dans la section Tables du lakehouse.
5.2. Ajouter des données à une table existante (append)
# Ajouter de nouvelles données à la table existante
df_nouvelles.write.mode("append").format("delta").saveAsTable("commandes")5.3. Sauvegarder avec partitionnement
# Partitionner par année et trimestre pour optimiser les requêtes
df_clean.write.mode("overwrite") \
.format("delta") \
.partitionBy("Annee", "Trimestre") \
.saveAsTable("commandes_partitionnees")Etape 6 : interroger la table Delta en SQL
Une fois la table créée, vous pouvez l’interroger directement en SQL Spark dans le notebook :
6.1. Requête de base
%%sql
SELECT * FROM commandes LIMIT 106.2. Agrégation : chiffre d’affaires par année
%%sql
SELECT
Annee,
COUNT(DISTINCT NumeroCommande) AS NbCommandes,
ROUND(SUM(MontantLigne), 2) AS ChiffreAffaires
FROM commandes
GROUP BY Annee
ORDER BY Annee6.3. Top 5 des produits les plus vendus
%%sql
SELECT
Produit,
SUM(Quantite) AS QuantiteTotale,
ROUND(SUM(MontantLigne), 2) AS CATotal
FROM commandes
GROUP BY Produit
ORDER BY QuantiteTotale DESC
LIMIT 56.4. Requête SQL dans du code PySpark
# Exécuter une requête SQL et récupérer le résultat dans un dataframe
df_resultats = spark.sql("""
SELECT
Annee,
Trimestre,
ROUND(SUM(MontantLigne), 2) AS CA
FROM commandes
GROUP BY Annee, Trimestre
ORDER BY Annee, Trimestre
""")
display(df_resultats)Etape 7 : créer une visualisation dans le notebook
Les notebooks Fabric permettent de créer des graphiques directement depuis les résultats :
import matplotlib.pyplot as plt
# Récupérer les données en Pandas pour la visualisation
df_ca = spark.sql("""
SELECT Annee, ROUND(SUM(MontantLigne), 2) AS CA
FROM commandes
GROUP BY Annee
ORDER BY Annee
""").toPandas()
# Créer un graphique en barres
plt.figure(figsize=(8, 5))
plt.bar(df_ca["Annee"].astype(str), df_ca["CA"], color="#2E75B6")
plt.xlabel("Année")
plt.ylabel("Chiffre d'affaires")
plt.title("Chiffre d'affaires par année")
plt.tight_layout()
plt.show()Cf. aussi : Ingestion avec un notebook
Exercices complémentaires
- Créer une table de dimension : à partir du dataframe, créez une table dim_clients contenant les colonnes NomClient et Email (sans doublons), puis sauvegardez-la en table Delta.
- Ajouter une colonne conditionnelle : ajoutez une colonne SegmentClient qui vaut "Premium" si le montant total des commandes du client dépasse 500 euros, et "Standard" sinon. Indice : utilisez une agrégation groupBy puis une jointure.
- Gérer les données incrémentales : simulez l’arrivée d’un nouveau fichier 2025.csv dans le dossier. Chargez uniquement ce fichier et ajoutez-le à la table existante avec mode("append").
- Interroger via le point de terminaison SQL : accédez au point de terminaison SQL analytics de votre lakehouse et écrivez une requête T-SQL équivalente à la requête du top 5 des produits (étape 6.3).
Sources :
- Source (Microsoft Learn) - Charger des données dans un lakehouse avec un notebook
- Source (Microsoft Learn) - Transformer et requêter avec Spark
- Source (Microsoft Learn) - Tutoriel lakehouse : préparer et transformer les données
- Source (Microsoft Learn) - Lire et écrire avec Pandas
- Source (Microsoft Learning) - Analyser des données avec Apache Spark dans Fabric