For AI agents: the complete documentation index is available at https://docs.ovhcloud.com/fr/llms.txt, the full documentation bundle is available at https://docs.ovhcloud.com/fr/llms-full.txt, and this page is available as Markdown at https://docs.ovhcloud.com/fr/guides/public-cloud/data-platform/tutorials-pyspark-cheat-sheet.md.

Aide-mémoire PySpark : fonctions essentielles et avancées sur l'OVHcloud Data Platform

Voir en Markdown

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

  1. 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.
  2. Lakehouse Manager : créez les tables correspondantes dans le Lakehouse Manager.
  3. DPE (Data Processing Engine) : assurez-vous que les données sont chargées dans ces tables au sein du DPE.
  4. 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

from forepaas.dwh import connect
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, unix_timestamp, hour, when, count, avg, udf, approx_count_distinct, collect_list, explode, regexp_replace, row_number
from pyspark.sql.types import FloatType
from pyspark.sql.window import Window
from forepaas.dwh.common import request as dwh_request, DwhRequestException
import io
import logging

# Configurer le logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

# Configurer les variables (DATASET, PROJECT_ID, YEAR, MONTH)
DATASET = "default_dataset"
PROJECT_ID = "PROJECT_ID" # Assurez-vous de remplacer ceci par votre PROJECT_ID réel
YEAR = "2025"
MONTH = "01"

# Initialiser la SparkSession
spark = SparkSession.builder.appName("PySpark_Advanced_Cheat_Sheet").getOrCreate()
logging.info(f"Spark Version: {spark.version}")

# Se connecter au Lakehouse
cn_prim = connect("dwh/default_dataset/")

# Charger les données
taxi_df = cn_prim.query(f"SELECT * FROM db_{PROJECT_ID}_{DATASET}.{DATASET}.yellow_tripdata_{YEAR}_{MONTH}")
zones_df = cn_prim.query(f"SELECT LocationID, Borough, Zone FROM db_{PROJECT_ID}_{DATASET}.{DATASET}.taxi_zone_lookup")
taxi_df.cache()
zones_df.cache()
logging.info(f"Taxi Records: {taxi_df.count()}")
logging.info(f"Zones Records: {zones_df.count()}")

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.
  • 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.
  • 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

# Nettoyer et transformer
cleaned_df = taxi_df \
    .filter(
        (col("tpep_pickup_datetime").isNotNull()) &
        (col("fare_amount") > 0) &
        (col("trip_distance") > 0)
    ) \
    .select(
        "tpep_pickup_datetime",
        "tpep_dropoff_datetime",
        "trip_distance",
        "fare_amount",
        col("pulocationid").cast("double").alias("pulocationid"),
        regexp_replace(col("store_and_fwd_flag"), "^[Yy]$", "Yes").alias("store_and_fwd_flag")
    ) \
    .withColumn(
        "trip_duration",
        unix_timestamp("tpep_dropoff_datetime") - unix_timestamp("tpep_pickup_datetime")
    ) \
    .withColumn(
        "pickup_hour",
        hour("tpep_pickup_datetime")
    ) \
    .withColumn(
        "trip_duration",
        when(col("trip_duration") > 3600, 3600).otherwise(col("trip_duration"))
    ) \
    .filter(col("trip_duration") >= 60)

logging.info(f"Cleaned Records: {cleaned_df.count()}")

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 dataType spécifié.
    • Bénéfice : garantit la compatibilité des types de données pour les calculs ou les traitements ultérieurs (par ex. convertir pulocationid en double).
  • 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 condition est vraie, la colonne prend value ; sinon, elle prend other_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).

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

# UDF pour l'efficacité tarifaire (tarif par minute)
def fare_efficiency(fare, duration):
    return fare / (duration / 60) if duration > 0 else 0.0

fare_efficiency_udf = udf(fare_efficiency, FloatType())

# Appliquer l'UDF
transformed_df = cleaned_df \
    .withColumn("fare_efficiency", fare_efficiency_udf(col("fare_amount"), col("trip_duration")))

transformed_df.show(5)

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() (depuis pyspark.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_udf nouvellement définie afin de créer la colonne fare_efficiency.

Résultat de ce bloc de code :

  • Une sortie DataFrame.show(5), affichant les 5 premières lignes de transformed_df, y compris la colonne fare_efficiency nouvellement ajoutée.
    • Exemple : si fare_amount vaut 15.0 et trip_duration vaut 600 secondes (10 minutes), fare_efficiency vaudra 1.5 ($/min).

Étape 4 : Window Functions avancées et jointures

Bloc de code

# Définir une window pour classer les trajets par tarif au sein d'un borough
window_spec = Window.partitionBy("pickup_borough").orderBy(col("fare_amount").desc())

# Joindre avec les zones et classer les trajets
joined_df = transformed_df \
    .join(
        zones_df,
        transformed_df.pulocationid == zones_df.LocationID,
        "left"
    ) \
    .withColumnRenamed("Borough", "pickup_borough") \
    .drop("LocationID") \
    .filter(col("pickup_borough").isNotNull()) \
    .withColumn("fare_rank", row_number().over(window_spec))

logging.info(f"Joined Records: {joined_df.count()}")
joined_df.filter(col("fare_rank") <= 3).show()

Fonctions utilisées dans ce bloc de code :

  • Window.partitionBy(*cols).orderBy(*cols) :

    • Quoi : définit une spécification de window. partitionBy divise les lignes en groupes, et orderBy dé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.
  • 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).
  • 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_rank est inférieur ou égal à 3, montrant les 3 meilleurs tarifs par borough.

Étape 5 : agréger et pivoter les données

Bloc de code

# Agréger : zones uniques et trajets par borough
agg_df = joined_df \
    .groupBy("pickup_borough") \
    .agg(
        approx_count_distinct("pulocationid").alias("unique_zones"),
        count("*").alias("num_trips")
    )

# Pivoter : tarif moyen par heure et par borough
pivot_df = joined_df \
    .groupBy("pickup_hour") \
    .pivot("pickup_borough") \
    .agg(avg("fare_amount")) \
    .orderBy("pickup_hour")

agg_df.show()
pivot_df.show()

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 countDistinct pour 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.
  • 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 avec pickup_borough, unique_zones (approximatif), et num_trips (par ex. Manhattan pourrait afficher ~60 zones uniques et ~2 millions de trajets).
  • pivot_df.show() : affiche une table pivotée montrant pickup_hour en lignes et pickup_borough en colonnes, avec la moyenne de fare_amount dans chaque cellule.

Étape 6 : collecter et exploser des listes

Bloc de code

# Collecter les zones par borough
list_df = joined_df \
    .groupBy("pickup_borough") \
    .agg(collect_list("Zone").alias("zones_list"))

# Exploser la liste des zones
exploded_df = list_df \
    .select("pickup_borough", explode(col("zones_list")).alias("zone"))

exploded_df.show(10)

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ée N lignes 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.

Résultat de ce bloc de code :

  • Une sortie DataFrame.show(10), affichant une table avec pickup_borough et une colonne zone explosé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

def create_table_from_this_dataframe(dataframe, dataset, table_name, bucket, source_bucket):
    logging.info(f"We will create a source (bucket) - {source_bucket} where we will store the new table - {table_name} - and automatically load it")

    # Création du bucket pour stocker la table
    cn_datastore = connect('data_store')
    logging.info(f"{cn_datastore.list()} - Before creating new bucket")
    cn_datastore.create_bucket(bucket)
    logging.info(f"{cn_datastore.list()} - After adding new bucket")
    cn_bucket = connect('data_store/' + bucket)

    get_dbs = dwh_request(f"v4/databases", method="GET")
    data = get_dbs.json()
    db_exist = next((item['_id'] for item in data if item.get('display_name') == source_bucket and item.get("package") == "data-store"), None)

    if db_exist is None:
        # Création de la source où nous utiliserons le nouveau bucket créé pour accéder à la table
        new_source = {"type":"protocol","package":"data-store","parameters":{"path":"","bucket":bucket},"default":False,"level":"source","display_name":source_bucket}
        new_source_bucket = dwh_request(f"v4/databases", method="POST", json=new_source)
        logging.info(f"New source added with the bucket: {bucket} - source name: {source_bucket}")
    else:
        logging.info("Source already exist")

    # Appel à l'API - Pour obtenir l'id du dataset
    get_database_id = dwh_request(f"v4/databases", method="GET")
    database_all = get_database_id.json()
    # Filtrage pour obtenir l'_id correspondant pour la base de données
    dataset_id = next((item['_id'] for item in database_all if item.get('name') == dataset), None)

    logging.info(f"dataset_id : {dataset_id}")

    # Convertir la table en Pandas et sérialiser en CSV dans BytesIO
    try:
        # Convertir en DataFrame Pandas
        table = dataframe.toPandas()

        # Créer le buffer BytesIO et écrire le CSV
        data = io.BytesIO()
        table.to_csv(data, index=False, encoding='utf-8')
        data.seek(0)  # Réinitialiser la position du buffer

        # Charger vers le bucket
        file_path = f"{table_name}.csv"
        etag = cn_bucket.put(file_path, data, data.getbuffer().nbytes)
        logging.info(f"DataFrame uploaded to bucket {bucket}/{file_path} with ETag: {etag}")

        # Vérifier le contenu du bucket
        files = cn_bucket.list()
        logging.info(f"Bucket contents: {files}")
    except Exception as e:
        logging.error(f"Failed to save to bucket: {e}")
        raise

    # Appel à l'API - Pour ajouter le fichier dans la source
    table_config = {"display_name":file_path,"progress":None,"physical_status":None,"parameters":{},"filename":file_path}
    res_table = dwh_request(f"v4/databases/{source_bucket}/tables/{file_path}", method="PUT", json=table_config)

    # Appel à l'API - Pour obtenir l'ID correspondant au file_path ajouté dans la source
    get_template_catalog_object = dwh_request(f"v4/tables", method="GET")
    data = get_template_catalog_object.json()
    file_path_source_id = next((item['_id'] for item in data if item.get('filename') == file_path), None)

    logging.info(f"file_path_source_id : {file_path_source_id}")

    # Au cas où la table existe déjà et que vous y avez apporté des modifications
    auto_build_table_DELETE = dwh_request(f"v4/logical/objects/{table_name}", method="DELETE") 

    # Appel à l'API - Pour lancer le build de la table sur le dataset correspondant et charger les données spécifiques
    config_build = {"database_id":dataset_id,"display_name":table_name,"name":table_name,"type":"prim","load_data":True,"build_table":True,"templated_from":"data_catalog","template_catalog_object":file_path_source_id}
    auto_build_table = dwh_request(f"v4/logical/objects", method="POST", json=config_build)

    logging.info(f"You can check the build of the table {table_name} on the Lakehouse Manager screen")

# Utiliser la fonction create_table_from_this_dataframe :
create_table_from_this_dataframe(pivot_df,"default_dataset","taxi_pivot_table","new_bucket","new_source_bucket")

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 io qui 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.
  • cn_datastore.create_bucket(bucket_name) (depuis forepaas.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) (depuis forepaas.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() (depuis forepaas.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) (depuis forepaas.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 connect de 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 :

  1. Création du bucket : un nouveau bucket (new_bucket) est créé dans le data store OVHcloud s'il n'existe pas déjà.
  2. 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.
  3. Sérialisation et chargement des données : la pivot_df est 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éé.
  4. 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.
Warning

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

# Vérifier que la table existe dans le Lakehouse
try:
    test_result = cn_prim.query(f"SELECT * FROM db_{PROJECT_ID}_{DATASET}.{DATASET}.taxi_pivot_table")
    logging.info(f"taxi_pivot_table Records: {test_result.count()}")
except Exception as e:
    logging.error(f"Failed to query table - please ensure it's the correct table")

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_table a bien été créée et est interrogeable dans le Lakehouse.
  • 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: 24 si 24 heures sont présentes).
  • En cas d'échec : un message d'erreur provenant de logging.error indiquant que la table n'a pas pu être interrogée.

Étape 9 : nettoyer et arrêter la session Spark

Bloc de code

spark.catalog.clearCache()
spark.stop()
logging.info("Spark cache cleared and Spark session stopped.")

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 à :

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.

Cette page vous a-t-elle aidé ?