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.

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=Trueretourne un DataFrame Fugueas_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"
)

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 :
- Conserver le style de programmation Pandas familier
- Éviter d'avoir à comprendre en détail les complexités du calcul distribué
- Migrer sans effort entre différents environnements de calcul
- 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 .