Aide-mémoire PySpark : fonctions essentielles et avancées sur l'OVHcloud Data Platform
Cet aide-mémoire fournit un guide concis des fonctions PySpark essentielles et avancées pour le traitement de données sur l'OVHcloud Data Platform
Objectif
Cet aide-mémoire fournit un guide concis des fonctions PySpark essentielles et avancées pour le traitement de données sur l'OVHcloud Data Platform. Il utilise les NYC Yellow Taxi Trip Records de janvier 2025 (yellow_tripdata_2025_01.parquet, ~3,5 millions d'enregistrements) et la Taxi Zone Lookup Table (taxi_zone_lookup.csv, 265 enregistrements) comme exemples pratiques.
Destiné aux utilisateurs intermédiaires, il couvre les fonctions de base (filter, select, groupBy, join, udf) et des fonctions avancées (window, pivot, approx_count_distinct, collect_list, explode, regexp_replace) pour des transformations et analyses complexes. Ces exemples illustrent des applications pratiques, faisant de ce document une référence polyvalente pour tout dataset.
Prérequis
Avant de commencer, assurez-vous de disposer de :
- Datasets : (
yellow_tripdata_2025_01.parquet,taxi_zone_lookup.csv) disponibles dans vos Connectors et accessibles dans le Lakehouse Manager. Vous pouvez les télécharger depuis le site officiel NYC TLC Trip Record Data. - Notebook : un notebook Jupyter compatible PySpark.
Instructions de configuration
- Connectors : créez une source nommée « NYC-taxi », chargez vos fichiers (
yellow_tripdata_2025_01.parquet,taxi_zone_lookup.csv), et extrayez leurs schémas. - Lakehouse Manager : créez les tables correspondantes dans le Lakehouse Manager.
- DPE (Data Processing Engine) : assurez-vous que les données sont chargées dans ces tables au sein du DPE.
- Notebook : démarrez un notebook Jupyter PySpark au sein de l'OVHcloud Data Platform.
Aide-mémoire des fonctions PySpark
Étape 1 : initialiser Spark et charger les données
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
SparkSession.builder.appName().getOrCreate():- Quoi : initialise la session PySpark et la DataFrame API, qui est le point d'entrée pour utiliser les fonctionnalités de Spark.
- Bénéfice : configure l'environnement pour le traitement de données distribué.
-
cn_prim = connect("dwh/default_dataset/"):- Quoi : établit une connexion vers votre dataset Lakehouse spécifié en utilisant la bibliothèque
forepaas.dwh. - Bénéfice : permet de requêter et d'interagir avec les tables stockées dans votre Lakehouse.
- Quoi : établit une connexion vers votre dataset Lakehouse spécifié en utilisant la bibliothèque
-
cn_prim.query(sql_query):- Quoi : exécute une requête SQL sur les tables accessibles via la connexion
cn_prim. Elle récupère les données sous forme de DataFrame Spark. - Bénéfice : fournit un moyen simple de charger des données depuis vos Connectors et votre Lakehouse dans une DataFrame Spark.
- Quoi : exécute une requête SQL sur les tables accessibles via la connexion
-
DataFrame.cache():- Quoi : marque la DataFrame pour qu'elle soit mise en cache en mémoire dès son premier calcul. Les opérations suivantes sur cette DataFrame liront depuis le cache.
- Bénéfice : accélère considérablement les opérations sur la même DataFrame, particulièrement utile pour les algorithmes itératifs ou les transformations multiples sur un grand dataset (~3,5 millions d'enregistrements).
-
DataFrame.count():- Quoi : déclenche le calcul et renvoie le nombre total de lignes de la DataFrame.
- Bénéfice : utilisé ici pour vérifier rapidement que les données ont été chargées et pour contrôler le nombre d'enregistrements.
Résultat de ce bloc de code :
- Version de Spark (par ex.
Spark Version: 3.4.1) Taxi Records: ~3,500,000(le nombre réel peut légèrement varier)Zones Records: 265
Étape 2 : nettoyer et transformer les données
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
DataFrame.filter(condition):- Quoi : filtre les lignes de la DataFrame selon une condition donnée, en renvoyant une nouvelle DataFrame ne contenant que les lignes qui satisfont la condition.
- Bénéfice : essentiel pour le nettoyage des données, en supprimant les enregistrements invalides ou non pertinents.
-
pyspark.sql.functions.col(column_name):- Quoi : référence une colonne dans une DataFrame, permettant d'y appliquer diverses transformations et opérations.
- Bénéfice : fournit un moyen de construire des expressions impliquant des colonnes de la DataFrame.
-
DataFrame.select(columns):- Quoi : projette un ensemble d'expressions (colonnes ou transformations basées sur des colonnes) et renvoie une nouvelle DataFrame contenant uniquement ces colonnes sélectionnées.
- Bénéfice : utile pour restreindre les colonnes et réduire la taille de la DataFrame en ne sélectionnant que les champs pertinents.
-
Column.cast(dataType):- Quoi : convertit le type de données d'une colonne vers le
dataTypespécifié. - Bénéfice : garantit la compatibilité des types de données pour les calculs ou les traitements ultérieurs (par ex. convertir
pulocationidendouble).
- Quoi : convertit le type de données d'une colonne vers le
-
pyspark.sql.functions.regexp_replace(column, pattern, replacement):- Quoi : remplace toutes les occurrences d'un motif de chaîne dans les données textuelles d'une colonne par une chaîne de remplacement spécifiée.
- Bénéfice : excellent pour la standardisation des données et le nettoyage des champs de type chaîne (par ex. convertir "Y" en "Yes").
-
DataFrame.withColumn(colName, col):- Quoi : renvoie une nouvelle DataFrame en ajoutant une nouvelle colonne ou en remplaçant une colonne existante par l'expression spécifiée.
- Bénéfice : essentiel pour le feature engineering et la création de nouvelles colonnes dérivées à la volée.
-
pyspark.sql.functions.unix_timestamp(timestamp_column):- Quoi : convertit une chaîne de timestamp ou une colonne de timestamp en timestamp Unix (secondes depuis le 1970-01-01 00:00:00 UTC).
- Bénéfice : facilite les calculs de temps numériques, comme le calcul des durées de trajet en soustrayant des timestamps.
-
pyspark.sql.functions.hour(timestamp_column):- Quoi : extrait la composante heure d'une colonne de timestamp.
- Bénéfice : utile pour l'analyse temporelle, permettant d'identifier des tendances ou motifs horaires dans les données.
-
pyspark.sql.functions.when(condition, value).otherwise(other_value):- Quoi : implémente une logique conditionnelle. Si la
conditionest vraie, la colonne prendvalue; sinon, elle prendother_value. Peut être chaîné. - Bénéfice : efficace pour gérer les valeurs aberrantes, appliquer des règles métier ou catégoriser des données selon des conditions spécifiques (par ex. plafonner
trip_duration).
- Quoi : implémente une logique conditionnelle. Si la
Résultat de ce bloc de code :
Cleaned Records: ~2,700,000 – ~2,800,000(le nombre réel peut varier selon la qualité des données).
Étape 3 : logique personnalisée avec UDF (User-Defined Function)
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
udf(func, returnType):- Quoi : enregistre une fonction Python comme User-Defined Function (UDF) dans PySpark. Cela permet d'appliquer une logique Python personnalisée aux colonnes d'une DataFrame Spark.
- Bénéfice : permet des calculs sur mesure et complexes qui pourraient ne pas être disponibles dans les fonctions natives de PySpark (par ex. calculer le « tarif par minute » avec une logique personnalisée).
-
FloatType()(depuispyspark.sql.types) :- Quoi : spécifie le type de données de retour de l'UDF comme un Float (nombre à virgule flottante en simple précision).
- Bénéfice : garantit que la sortie de votre fonction personnalisée est correctement typée dans la DataFrame Spark.
-
DataFrame.withColumn(colName, col):- Quoi : (réutilisé depuis l'étape 2) ajoute une nouvelle colonne ou remplace une colonne existante selon le résultat d'une expression.
- Bénéfice : utilisé ici pour appliquer la
fare_efficiency_udfnouvellement définie afin de créer la colonnefare_efficiency.
Résultat de ce bloc de code :
- Une sortie
DataFrame.show(5), affichant les 5 premières lignes detransformed_df, y compris la colonnefare_efficiencynouvellement ajoutée.- Exemple : si
fare_amountvaut 15.0 ettrip_durationvaut 600 secondes (10 minutes),fare_efficiencyvaudra 1.5 ($/min).
- Exemple : si
Étape 4 : Window Functions avancées et jointures
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
Window.partitionBy(*cols).orderBy(*cols):- Quoi : définit une spécification de window.
partitionBydivise les lignes en groupes, etorderBydéfinit l'ordre logique des lignes au sein de chaque partition. - Bénéfice : essentiel pour permettre des opérations analytiques avancées (comme le classement, lead/lag, les sommes cumulatives) qui s'appliquent à un sous-ensemble défini de lignes.
- Quoi : définit une spécification de window.
-
pyspark.sql.functions.row_number():- Quoi : une window function qui attribue un numéro séquentiel unique à chaque ligne au sein de sa partition, selon l'ordre défini dans la spécification de window.
- Bénéfice : parfait pour classer des enregistrements (par ex. identifier les N meilleurs enregistrements selon une métrique, comme les tarifs les plus élevés).
-
DataFrame.join(other_df, on=None, how=None):- Quoi : combine deux DataFrames selon une condition de jointure (
on) et un type de jointure (how, par ex. "inner", "left", "right") spécifiés. - Bénéfice : enrichit les données en rassemblant des informations liées provenant de différentes sources (par ex. joindre les données de trajet de taxi avec les données de lookup des zones).
- Quoi : combine deux DataFrames selon une condition de jointure (
-
DataFrame.withColumnRenamed(existing, new):- Quoi : renvoie une nouvelle DataFrame en renommant une colonne existante.
- Bénéfice : aide à clarifier le schéma et améliore la lisibilité, en particulier après des jointures où les noms de colonnes peuvent être ambigus.
-
DataFrame.drop(*cols):- Quoi : renvoie une nouvelle DataFrame avec les colonnes spécifiées supprimées.
- Bénéfice : aide à gérer la taille et la complexité de la DataFrame en supprimant les colonnes inutiles, économisant de la mémoire.
Résultat de ce bloc de code :
Joined Records: ~2,600,000 – ~2,700,000(le nombre réel peut varier).- Une sortie
DataFrame.show(), affichant les lignes oùfare_rankest inférieur ou égal à 3, montrant les 3 meilleurs tarifs par borough.
Étape 5 : agréger et pivoter les données
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
DataFrame.groupBy(*cols):- Quoi : regroupe la DataFrame selon une ou plusieurs colonnes spécifiées, en préparation des calculs d'agrégation.
- Bénéfice : permet la synthèse et l'analyse des données selon des catégories ou dimensions distinctes.
-
DataFrame.agg(*exprs):- Quoi : applique des fonctions d'agrégation aux données regroupées, en calculant des statistiques de synthèse.
- Bénéfice : utilisé pour calculer des métriques comme des comptages, sommes, moyennes, etc., pour chaque groupe.
-
pyspark.sql.functions.approx_count_distinct(column):- Quoi : renvoie un comptage approximatif d'éléments distincts au sein d'un groupe. Utilise l'algorithme HyperLogLog++.
- Bénéfice : nettement plus rapide et plus efficace en mémoire que
countDistinctpour de très grands datasets lorsqu'un comptage exact n'est pas strictement nécessaire.
-
pyspark.sql.functions.count(column):- Quoi : compte le nombre de valeurs non nulles dans une colonne ou, avec
count("*"), compte toutes les lignes d'un groupe. - Bénéfice : comptabilise la taille de chaque groupe agrégé ou les occurrences de valeurs spécifiques.
- Quoi : compte le nombre de valeurs non nulles dans une colonne ou, avec
-
DataFrame.pivot(pivot_column):- Quoi : pivote une DataFrame, transformant les valeurs uniques d'une colonne spécifiée en nouvelles colonnes. Nécessite une agrégation ultérieure.
- Bénéfice : crée des tables larges, souvent plus adaptées au reporting et à l'analyse transversale, permettant de comparer directement des valeurs entre catégories.
-
pyspark.sql.functions.avg(column):- Quoi : calcule la valeur moyenne d'une colonne numérique.
- Bénéfice : fournit une mesure de tendance centrale pour les données quantitatives au sein de chaque groupe.
-
DataFrame.orderBy(*cols, ascending=True):- Quoi : trie les lignes de la DataFrame selon une ou plusieurs colonnes, par ordre croissant ou décroissant.
- Bénéfice : organise la sortie pour une meilleure lisibilité et pour présenter les données dans un ordre logique.
Résultat de ce bloc de code :
agg_df.show(): affiche une table avecpickup_borough,unique_zones(approximatif), etnum_trips(par ex. Manhattan pourrait afficher ~60 zones uniques et ~2 millions de trajets).pivot_df.show(): affiche une table pivotée montrantpickup_houren lignes etpickup_boroughen colonnes, avec la moyenne defare_amountdans chaque cellule.
Étape 6 : collecter et exploser des listes
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
pyspark.sql.functions.collect_list(column):- Quoi : une fonction d'agrégation qui rassemble toutes les valeurs non nulles d'une colonne spécifiée au sein de chaque groupe dans une liste Python.
- Bénéfice : utile pour créer des structures de type tableau où chaque élément correspond à un enregistrement du groupe d'origine.
-
pyspark.sql.functions.explode(array_column):- Quoi : transforme une colonne contenant des tableaux (listes) ou des maps en lignes individuelles pour chaque élément du tableau/map. Si un tableau contient
Néléments, cela créeNlignes pour cette ligne d'origine. - Bénéfice : aplatit les structures de données imbriquées, permettant de traiter ou de visualiser des éléments individuels comme des enregistrements séparés.
- Quoi : transforme une colonne contenant des tableaux (listes) ou des maps en lignes individuelles pour chaque élément du tableau/map. Si un tableau contient
Résultat de ce bloc de code :
- Une sortie
DataFrame.show(10), affichant une table avecpickup_boroughet une colonnezoneexplosée, où chaque zone distincte d'un borough obtient sa propre ligne (par ex. pour Manhattan, vous verriez plusieurs lignes comme "Manhattan | Midtown", "Manhattan | Upper East Side", etc.).
Étape 7 : sauvegarder la DataFrame dans un Bucket et créer une table (avancé)
Bloc de code
Fonctions utilisées dans ce bloc de code (et au sein de create_table_from_this_dataframe) :
-
DataFrame.toPandas():- Quoi : convertit une DataFrame Spark en DataFrame Pandas. Cela rassemble toutes les données distribuées sur le nœud driver.
- Bénéfice : permet d'utiliser les fonctions spécifiques à Pandas pour la manipulation locale des données et la sérialisation de fichiers (par ex.
to_csv). - Avertissement : à utiliser avec prudence sur de très grands datasets, car cela peut provoquer des erreurs de mémoire insuffisante sur le driver.
-
Pandas_DataFrame.to_csv(path_or_buffer, index=False, encoding='utf-8'):- Quoi : écrit la DataFrame Pandas dans un fichier de valeurs séparées par des virgules (CSV).
- Bénéfice : sérialise les données dans un format texte standard adapté au stockage et à la récupération.
-
io.BytesIO():- Quoi : une classe du module Python
ioqui crée un flux binaire en mémoire, se comportant comme un objet fichier. - Bénéfice : permet d'écrire et de lire des octets comme si vous interagissiez avec un fichier physique, ce qui est utile pour le transfert direct de données vers des services sans enregistrement sur disque.
- Quoi : une classe du module Python
-
cn_datastore.create_bucket(bucket_name)(depuisforepaas.dwh) :- Quoi : crée un nouveau bucket de stockage au sein du data store de l'OVHcloud Data Platform.
- Bénéfice : fournit un emplacement dédié pour stocker des fichiers, y compris des données intermédiaires ou finales traitées.
-
cn_bucket.put(file_path, data, size)(depuisforepaas.dwh) :- Quoi : charge des données (généralement depuis un buffer BytesIO) vers un chemin spécifié au sein d'un bucket connecté.
- Bénéfice : persiste vos données traitées (par ex. le CSV issu de
pivot_df) dans le stockage cloud OVHcloud.
-
cn_bucket.list()(depuisforepaas.dwh) :- Quoi : récupère une liste des fichiers et sous-répertoires au sein d'un bucket connecté.
- Bénéfice : utilisé pour vérifier que les fichiers ont bien été chargés dans le bucket.
-
dwh_request(path, method, json)(depuisforepaas.dwh.common) :- Quoi : une fonction utilitaire pour effectuer des appels API HTTP directs vers les services backend de l'OVHcloud Data Platform.
- Bénéfice : cette fonction offre un contrôle granulaire pour automatiser des tâches comme la création de sources de données, l'enregistrement de fichiers dans les Connectors, et le déclenchement de builds de table dans le Lakehouse Manager de manière programmatique, ce qui n'est pas toujours exposé via l'objet
connectde plus haut niveau.
Explication du processus de cette étape :
Cette étape avancée montre comment persister une DataFrame Spark traitée (pivot_df) dans une nouvelle table du Lakehouse de l'OVHcloud Data Platform. Elle s'appuie sur une combinaison de PySpark, Pandas et d'appels API directs :
- Création du bucket : un nouveau bucket (
new_bucket) est créé dans le data store OVHcloud s'il n'existe pas déjà. - Création de la source : une nouvelle source de données (
new_source_bucket) est configurée et reliée à ce bucket, permettant aux Connectors de découvrir les fichiers qu'il contient. - Sérialisation et chargement des données : la
pivot_dfest convertie en DataFrame Pandas, puis sérialisée au format CSV dans un buffer en mémoire (io.BytesIO). Ces données CSV sont ensuite chargées vers le bucket nouvellement créé. - Enregistrement et build de la table : via des appels API directs (
dwh_request), le fichier CSV chargé est enregistré comme « table physique » dans les Connectors, et une « table logique » (taxi_pivot_table) est créée dans le Lakehouse Manager, déclenchant un chargement des données depuis la source enregistrée.
Avertissement : il s'agit d'une solution temporaire pour sauvegarder et créer des tables, reposant sur des interactions API directes. Le SDK de l'OVHcloud Data Platform devrait être mis à jour pour inclure des capacités de transformation SQL et de stockage direct plus natives et simplifiées à l'avenir, réduisant le besoin d'appels API manuels et de conversions Pandas.
Étape 8 : vérifier l'existence de la table
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
cn_prim.query(sql_query):- Quoi : (réutilisé depuis l'étape 1) exécute une requête SQL sur les tables de votre Lakehouse via la connexion
cn_primétablie, en récupérant le résultat sous forme de DataFrame Spark. - Bénéfice : utilisé ici pour confirmer que
taxi_pivot_tablea bien été créée et est interrogeable dans le Lakehouse.
- Quoi : (réutilisé depuis l'étape 1) exécute une requête SQL sur les tables de votre Lakehouse via la connexion
-
DataFrame.count():- Quoi : (réutilisé depuis l'étape 1) renvoie le nombre total de lignes de la DataFrame interrogée.
- Bénéfice : confirme que la table nouvellement créée contient les données et enregistrements attendus.
Résultat de ce bloc de code :
- En cas de succès :
taxi_pivot_table Records: [Nombre de lignes dans pivot_df](par ex.taxi_pivot_table Records: 24si 24 heures sont présentes). - En cas d'échec : un message d'erreur provenant de
logging.errorindiquant que la table n'a pas pu être interrogée.
Étape 9 : nettoyer et arrêter la session Spark
Bloc de code
Fonctions utilisées dans ce bloc de code :
-
SparkSession.catalog.clearCache():- Quoi : vide le cache interne de Spark, libérant la mémoire occupée par les DataFrames et RDD mis en cache.
- Bénéfice : essentiel pour libérer des ressources, en particulier après des opérations complexes ou lorsque vous n'avez plus besoin des données mises en cache. Cela aide à prévenir les problèmes de mémoire insuffisante dans les applications longue durée ou les sessions interactives.
-
SparkSession.stop():- Quoi : termine la SparkSession, arrêtant le SparkContext et libérant toutes les ressources associées.
- Bénéfice : garantit un arrêt propre de votre application Spark, libérant les ressources du cluster. Il est crucial d'appeler cette fonction à la fin de votre application Spark pour éviter les fuites de ressources.
Résultat de ce bloc de code :
INFO - Spark cache cleared and Spark session stopped.- Ce message confirme que le cache a été vidé et que la session Spark a été correctement terminée. Vous verrez généralement une sortie supplémentaire de l'environnement du notebook indiquant que le kernel est inactif ou s'est arrêté.
Conclusion et étapes suivantes
Cet aide-mémoire PySpark a fourni un aperçu pratique des fonctions essentielles et avancées pour le traitement de données sur l'OVHcloud Data Platform. Vous avez vu comment initialiser Spark, charger et nettoyer des données, appliquer une logique personnalisée avec des UDF, effectuer des agrégations complexes avec des window functions et des pivots, et même gérer la persistance des données vers votre Lakehouse.
Ce guide sert de référence pratique pour accélérer votre développement PySpark. Pour poursuivre votre parcours et appliquer ces compétences à des scénarios réels, nous vous encourageons à :
- Explorer les tutoriels : plongez dans les tutoriels complets d'analyse du dataset NYC Taxi pour construire des pipelines de données de bout en bout et des modèles de machine learning.
- Approfondir vos connaissances : consultez la documentation officielle de PySpark pour des détails approfondis sur toutes les fonctions.
Bon code et bonne analyse de données !
Aller plus loin
Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.
Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.
Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.
Rejoignez notre communauté d'utilisateurs.