Parallélisation efficace de code Pandas avec la fonction transform de Fugue

Pourquoi paralléliser le code Pandas ?

Lors du traitement de grands ensembles de données, un processus Pandas unique rencontre souvent des goulets d'étranglement de performance :

  • Dépassement de mémoire lorsque le volume de données excède la capacité RAM
  • Calculs complexes trop longs affectant la productivité
  • Impossibilité d'exploiter pleinement les CPU multi-cœurs ou les ressources de cluster distribué

La fonction transform de Fugue a été conçue précisément pour résoudre ces problèmes, vous permettant de :

  • Conserver votre code Pandas existant
  • Basculer sans effort entre différents moteurs d'exécution
  • Contrôler flexiblement les stratégies de partitionnement
  • Mettre en œuvre facilement des calculs distribués

Principe de fonctionnement de la fonction transform de Fugue

L'idée fondamentale de la fonction transform est de diviser un grand ensemble de données en plusieurs partitions logiques, puis d'appliquer votre fonction Pandas à chaque partition, et enfin de fusionner les résultats. Ce processus est entièrement géré automatiquement par Fugue, sans nécessiter une manipulation manuelle des détails distribués.

Architecture Fugue montrant comment la fonction transform distribue les tâches de calcul sur différents moteurs d'exécution

L'interface Transformer de Fugue définit les méthodes centrales de la logique de transformation, dont la plus importante est la méthode transform, qui reçoit un DataFrame local et retourne le résultat transformé :

def transform(self, df: LocalDataFrame) -> LocalDataFrame:
    """Logique de transformation d'un dataframe local à un autre."""
    raise NotImplementedError

Premiers pas : Utilisation de base

Utiliser la fonction transform de Fugue est très simple, il suffit de trois étapes pour paralléliser votre code Pandas :

1. Définir votre fonction Pandas

Commencez par écrire une fonction Pandas standard, par exemple pour trier par la colonne "valeur" et prendre la première ligne de chaque groupe :

def trier_premier(df: pd.DataFrame) -> pd.DataFrame:
    return df.sort_values("valeur").head(1)

2. Utiliser la fonction transform pour l'exécution

Ensuite, utilisez directement la fonction transform de Fugue pour appeler cette fonction :

from fugue import transform
import pandas as pd

# Création de données d'exemple
donnees = pd.DataFrame([[1, 15], [0, 5], [1, 8], [0, 25]], columns=["id", "valeur"])

# Transformation de base
resultat = transform(donnees, trier_premier, schema="*")
print(resultat.values.tolist())  # Affiche: [[0, 5]]

3. Ajouter un partitionnement pour la parallélisation

Pour activer la parallélisation, ajoutez simplement le paramètre partition pour spécifier la méthode de partitionnement :

# Traitement parallèle par partition selon la colonne "id"
resultat = transform(donnees, trier_premier, partition=dict(by=["id"]))
print(sorted(resultat.values.tolist(), key=lambda x: x[0]))  # Affiche: [[0, 5], [1, 8]]

Fonctionnalités avancées et paramètres

La fonction transform de Fugue offre de nombreuses options de paramètres pour contrôler flexiblement le processus de parallélisation :

Spécification du mode de sortie

Contrôlez le type de sortie avec le paramètre as_fugue :

  • as_fugue=True retourne un DataFrame Fugue
  • as_fugue=False (défaut) retourne un DataFrame Pandas
# Retourner un DataFrame Fugue
df_fugue = transform(donnees, trier_premier, partition=dict(by=["id"]), as_fugue=True)

Persistance des données

Utilisez persist=True pour persister les résultats intermédiaires et éviter les calculs répétitifs :

resultat = transform(donnees, trier_premier, persist=True)

Sauvegarde des résultats dans un fichier

Sauvegardez directement les résultats dans un fichier avec le paramètre save_path :

transform(
    donnees, 
    trier_premier, 
    save_path="resultat.parquet",
    engine_conf={FUGUE_CONF_WORKFLOW_CHECKPOINT_PATH: "/chemin/vers/checkpoint"}
)

Étude de cas pratique : Traitement distribué de données

Supposons que vous disposiez d'un grand ensemble de données de ventes et que vous deviez calculer le total des ventes mensuelles par catégorie de produit pour chaque région. Avec la focntion transform de Fugue, vous pouvez facilement paralléliser cette tâche :

import pandas as pd
from fugue import transform

def ventes_mensuelles_par_categorie(df: pd.DataFrame) -> pd.DataFrame:
    # Calculer les ventes par mois et catégorie de produit
    df['mois'] = df['date'].dt.to_period('M')
    return df.groupby(['region', 'mois', 'categorie'])['montant'].sum().reset_index()

# Lire le grand ensemble de données
grand_df = pd.read_parquet("donnees_ventes importantes.parquet")

# Traitement parallèle par région
resultat = transform(
    grand_df, 
    ventes_mensuelles_par_categorie, 
    schema="region:str,mois:str,categorie:str,montant:float",
    partition=dict(by=["region"])
)

# Sauvegarder les résultats
resultat.to_parquet("rapport_ventes_mensuelles.parquet")

Extension à différents moteurs d'exécution

La puissance de Fugue réside dans sa capacité à basculer le calcul vers différents moteurs d'exécution sans modifier le code. Par exemple, pour exécuter la tâche précédente avec Spark, il suffit d'ajouter le paramètre engine :

# Exécution avec Spark
resultat = transform(
    grand_df, 
    ventes_mensuelles_par_categorie, 
    schema="region:str,mois:str,categorie:str,montant:float",
    partition=dict(by=["region"]),
    engine="spark"
)

Architecture d'extension Fugue montrant comment l'interface unifiée prend en charge plusieurs moteurs d'exécution

Ressources d'apprentissage

  • Documentation officielle : docs/
  • Code d'implémentation principal : fugue/extensions/transformer/transformer.py
  • Cas de test : tests/fugue/test_interfaceless.py

Avantages clés

La fonction transform de Fugue offre aux utilisateurs de Pandas une voie simple et efficace pour la parallélisation, leur permettant de :

  1. Conserver le style de programmation Pandas familier
  2. Éviter d'avoir à comprendre en détail les complexités du calcul distribué
  3. Migrer sans effort entre différents environnements de calcul
  4. Exploiter pleinement les ressources de calcul pour améliorer l'efficacité

Pour commencer à utiliser Fugue, clonez simplement le dépôt et installez-le :

git clone https://gitcode.com/gh_mirrors/fu/fugue
cd fugue
pip install .

Étiquettes: Fugue Pandas Parallélisation Spark Dask

Publié le 6 septembre à 23h45