Ingestion avec un notebook
Depuis un espace de travail ou dans un lakehouse, créer un notebook.
Récapitulatif des commandes clés
| Opération | Code PySpark |
| Lire un CSV | spark.read.format("csv").option("header", True).load("Files/fichier.csv") |
| Lire plusieurs CSV | spark.read.format("csv").schema(mon_schema).load("Files/dossier/*.csv") |
| Afficher les données | display(df) ou df.show() |
| Voir le schéma | df.printSchema() |
| Compter les lignes | df.count() |
| Supprimer les doublons | df.dropDuplicates() |
| Supprimer les nulls | df.dropna(subset=["colonne"]) |
| Ajouter une colonne | df.withColumn("nom", expression) |
| Renommer une colonne | df.withColumnRenamed("ancien", "nouveau") |
| Filtrer | df.filter(col("colonne") > valeur) |
| Agréger | df.groupBy("colonne").agg(sum("valeur")) |
| Sauver en Delta | df.write.mode("overwrite").format("delta").saveAsTable("table") |
| Ajouter à une table | df.write.mode("append").format("delta").saveAsTable("table") |
| Requête SQL | spark.sql("SELECT * FROM table") |
| SQL dans une cellule | %%sql en première ligne de la cellule |
Ingérer plusieurs tables de plusieurs lakehouse
Préparer les raccourcis
Dans un lakehouse, créer des raccourcis vers les schémas des autres lakehouse :
- Obtenir des données > Nouveau raccourci
- Sélectionner OneLake puis le lakehouse et Suivant
- Dans Tables, cocher dbo (ou le schéma concerné), puis Suivant
- Dans la dernière étape, cliquer sur Crayon dans Actions :
- Cliquer sur Créer.
Préparer le fichier des noms de tables
- Créer un fichier CSV NomTable.csv, écrire NomTable en première ligne, puis les noms des tables sut chaque ligne suivante.
- Obtenir des données > Charger des fichiers
- Cliquer sur le dossier pour afficher la boite Ouvrir
- Sélectionner le fichier CSV et Charger. Fermer le volet. Le fichier est dans le dossier Files.
Compléter le notebook
- Ouvrir le notebook > Nouveau notebook.
- Écrire le code suivant :
| <div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>from pyspark.sql.functions import lit</span></div> | |
| <div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span># Définition des clients et de leur lakehouse source</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>clients = [</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> {"client_name": "Client1", "dbo": "dboClient1"},</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> {"client_name": "Client2", "dbo": "dboClient2"}</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>]</span></div> | Remplacer Client1 pour le vrai nom du client, qui sera placé dans une nouvelle colonne.Remplacer dboClient1 par le nom du schéma lié. |
| <div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span># Lecture du CSV depuis la section Files du lakehouse</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>liste_tables_df = spark.read \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .option("header", "true") \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .option("delimiter", ";") \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .csv("Files/ListeTable.csv")</span></div> | Remplacer Files/ListeTable.csv par votre “vrai” chemin si besoin.Vérifier aussi la casse. |
| <div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>tables_erp = [</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> row["NomTable"]</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> for row in liste_tables_df.collect()</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> if row["NomTable"] is not None</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>]</span></div> | “NomTable” est indiqué sur la première ligne du fichier CSV |
| <div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span># Fusion pour chaque table</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span>for table_name in tables_erp:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> print(f"Traitement de : {table_name}")</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> dfs = []</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><br></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> for client in clients:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> try:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> df = spark.table(f"{client['dbo']}.{table_name}")</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> df = df.withColumn("client_name", lit(client["client_name"]))</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> dfs.append(df)</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> except Exception as e:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> print(f" Table {table_name} absente dans {client['lakehouse']} : {e}")</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> continue</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><br></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> if not dfs:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> print(f" Aucune source disponible pour {table_name}, table ignorée.")</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> continue</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><br></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> # df_final = dfs[0]</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> for df in dfs[0:]:</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> df_final = df_final.unionByName(df, allowMissingColumns=True)</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><br></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> df_final.write \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .format("delta") \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .mode("overwrite") \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .option("mergeSchema", "true") \</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> .saveAsTable(f"{table_name}_unified")</span></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><br></div><div style="text-align: left; margin-top: 0.5em; margin-bottom: 0.5em;"><span> print(f" Écrite dans Silver : {table_name}_unified")</span></div> | Boucle sur le liste des tables.Sous-boucle sur la liste des clients.Si une table n’existe pas dans le |
Enfin, exécuter le notebook
