---
title: "Aide-mémoire PySpark : fonctions essentielles et avancées sur l'OVHcloud Data Platform"
description: "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"
url: https://docs.ovhcloud.com/fr/guides/public-cloud/data-platform/tutorials-pyspark-cheat-sheet
lang: fr
lastUpdated: 2026-09-14
---
> 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.

# Aide-mémoire PySpark : fonctions essentielles et avancé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](https://www.nyc.gov/site/tlc/about/tlc-trip-record-data.page).
- **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**

```python
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**

```python
# 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**

```python
# 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**

```python
# 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**

```python
# 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**

```python
# 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**

```python
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**

```python
# 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**

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

- **Explorer les tutoriels** : plongez dans les [tutoriels complets d'analyse du dataset NYC Taxi](https://docs.ovhcloud.com/fr/guides/public-cloud/data-platform/tutorials-pyspark.md) 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](https://spark.apache.org/docs/latest/api/python/) 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](https://www.ovhcloud.com/fr/professional-services/) 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](https://discord.gg/ovhcloud) dédié.

Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre [centre d'aide](https://help.ovhcloud.com/csm?id=csm_get_help).

Rejoignez notre [communauté d'utilisateurs](https://community.ovhcloud.com/).
