Accéder au contenu principal

Maîtriser groupBy de PySpark pour une agrégation de données à l’échelle

Découvrez la méthode groupBy de PySpark, qui permet aux professionnels des données d’appliquer des fonctions d’agrégation à leurs données. Une façon puissante de partitionner et synthétiser rapidement vos grands jeux de données en tirant parti des techniques avancées de Spark.
Actualisé 19 sept. 2026  · 15 min lire

Explorer avec l’IA

ChatGPTClaudePerplexity

Dans l’univers du big data, on consulte rarement des données brutes. Pour les rendre exploitables, le regroupement et l’agrégation sont des opérations courantes et puissantes. En environnement distribué comme Apache Spark et Python, la fonction groupBy de PySpark nous permet de regrouper et de synthétiser nos données distribuées. Elle reflète le fonctionnement de la clause GROUP BY en SQL, tout en étant conçue pour traiter efficacement des jeux de données massifs en mode distribué.

La méthode groupBy de PySpark permet de partitionner les données selon différentes colonnes, puis d’agréger ces groupes avec des mesures telles que sommes, moyennes, etc. Pour optimiser cela, PySpark suit le paradigme split-apply-combine :

  1. Scinder les données en groupes selon des critères,
  2. Appliquer une logique d’agrégation ou de transformation à chaque groupe,
  3. Combiner les résultats dans un nouveau DataFrame.

Graphique montrant le fonctionnement de groupBy de PySpark : partition par clé, application d’une somme par clé, puis combinaison des données agrégées

Exemple d’application de Sum avec groupBy() de PySpark

PySpark exploite l’évaluation paresseuse : les opérations comme groupBy ne sont réellement calculées que lorsqu’une action (comme show() ou collect()) est appelée. PySpark construit d’abord un DAG (graphe acyclique orienté) initial, ensuite optimisé pour les performances. Cela aide Spark à optimiser le plan d’exécution avant de lancer le calcul. 

Si vous n’avez pas encore travaillé avec PySpark, nous vous recommandons vivement le cours Big Data Fundamentals with PySpark.

Créer un DataFrame PySpark

Avant tout, il nous faut un DataFrame PySpark. Besoin d’un rappel sur les commandes PySpark ? Consultez cet excellent aide-mémoire PySpark DataFrame

Passons en revue quelques étapes et commandes clés pour bien démarrer avec PySpark. Avant de commencer, assurez-vous d’avoir installé PySpark dans votre environnement Python, ainsi que Java et le JDK Java. Si besoin, suivez ce guide pour démarrer avec PySpark.

Initialiser SparkSession

Première étape : lancer PySpark !

from pyspark.sql import SparkSession

# Start a SparkSession
spark = SparkSession.builder \
    .appName("GroupByExample") \
    .getOrCreate()

DataFrame d’exemple pour GroupBy

Créons un DataFrame d’exemple sur lequel travailler.

from pyspark.sql import Row

data = [
    Row(department="Sales", employee="Alice", salary=5000),
    Row(department="Sales", employee="Bob", salary=4800),
    Row(department="HR", employee="Carol", salary=4000),
    Row(department="HR", employee="David", salary=3900),
    Row(department="IT", employee="Eve", salary=6000)
]

df = spark.createDataFrame(data)
df.show()

Résultat :

	+----------+--------+------+
	|department|employee|salary|
	+----------+--------+------+
	|     Sales|   Alice|  5000|
	|     Sales|     Bob|  4800|
	|        HR|   Carol|  4000|
	|        HR|   David|  3900|
	|        IT|     Eve|  6000|
	+----------+--------+------+

Qu’est-ce que GroupBy en PySpark ?

La transformation groupBy de PySpark permet de scinder les données en groupes selon une ou plusieurs colonnes, puis d’appliquer des agrégations ou transformations indépendantes.

Syntaxe de base et exemples

L’utilisation de groupBy est simple : appelez-la sur le DataFrame souhaité. Elle s’applique à tout type de colonne pertinent pour la partition. Par exemple, on évitera généralement les flottants.

grouped = df.groupBy("department")

Cela crée un objet GroupedData, qui n’est pas un DataFrame en soi. Il indique à Spark comment commencer à partitionner les données. Une fois cet objet créé, vous pouvez appliquer des agrégations comme .count() ou .sum() pour obtenir un nouveau DataFrame.

grouped.count().show()

Paramètres et valeurs de retour

Le seul paramètre de la méthode est *cols, qui accepte des noms de colonnes, des expressions de colonnes, des indices (int) ou une liste de colonnes. Tant que la colonne sert de clé de partition, vous pouvez l’utiliser. La méthode renvoie toujours un objet GroupedData.

# groupement par nom de colonne comme ci-dessus
	df.groupBy("department")

# groupement par expression de colonne
	df.groupBy(df.department)

# groupement par indice de colonne
	df.groupBy(1)

# groupement par liste de colonnes (vous pouvez combiner les méthodes)
	df.groupBy(["department", 2])

GroupBy sur une ou plusieurs colonnes

Comme montré ci-dessus, vous pouvez regrouper par une ou plusieurs colonnes. Une seule colonne convient lorsque vous vous intéressez à un seul axe (départements, années). Plusieurs colonnes apportent plus de granularité : commerciaux par département, mois au sein d’une année, etc.

# groupement par une seule colonne
	df.groupBy("department").sum(“salary”).show()

# groupement par plusieurs colonnes
	df.groupBy(["department", 2]).sum(“salary”).show()

Vous l’aurez noté, groupBy fonctionne de façon très similaire à GROUP BY en SQL. Nous verrons plus loin comment utiliser aussi ce langage dans PySpark.

Fonctions et techniques d’agrégation

PySpark propose un large éventail de méthodes d’agrégation intégrées pour traiter des données groupées. Si vous connaissez SQL, vous retrouverez des fonctions telles que count(), sum(), avg(), etc. Consultez notre guide des fonctions d’agrégation en SQL pour un rappel. 

Fonctions d’agrégation intégrées

Passons en revue les principales fonctions :

  • count() : compte le nombre d’enregistrements dans la partition
  • sum() : calcule la somme des valeurs numériques
  • avg() : calcule la moyenne des valeurs numériques
  • min() : renvoie la plus petite valeur de la partition
  • max() : renvoie la plus grande valeur de la partition

En complément, vous pouvez utiliser la méthode .alias() pour renommer les colonnes et rendre les résultats plus lisibles. Exemple ci-dessous.

Agrégations multiples avec agg()

Plutôt que d’exécuter chaque agrégation séparément, vous pouvez en effectuer plusieurs en une seule fois grâce à agg(). Chaque agrégation crée une nouvelle colonne dans votre DataFrame. Cela évite des appels multiples à groupBy et améliore les performances. Il n’est pas nécessaire d’agréger sur la même colonne : vous pouvez définir une colonne différente pour chaque fonction d’agrégation.

# Après avoir démarré votre session
from pyspark.sql import functions as sf

df.groupBy("department").agg(
    sf.count("employee").alias("employee_count"), 
    sf.avg("salary").alias("avg_salary"),
    sf.max("salary").alias("max_salary")
).show()

Résultat :

+----------+--------------+----------+----------+
|department|employee_count|avg_salary|max_salary|
+----------+--------------+----------+----------+
|     Sales|             2|    4900.0|      5000|
|        HR|             2|    3950.0|      4000|
|        IT|             1|    6000.0|      6000|
+----------+--------------+----------+----------+

Modèles d’agrégation avancés

L’un des atouts de PySpark est la prise en charge de modèles d’agrégation avancés : pivot, rollup et data cubes. Vous pouvez aussi créer des grouping sets. 

Pivot

Comme dans un tableau croisé dynamique Excel, vous pouvez pivoter vos données sur différentes colonnes. Ici, on groupe par department et on pivote sur employee pour voir le salaire total de chaque employé. 

Chaque ligne représente le département, et chaque colonne correspond à un employé du département. Imaginez la puissance de ce procédé avec des données annuelles pour comparer la performance d’un département d’une année sur l’autre.

df.groupBy("department").pivot("employee").sum("salary").show()

Rollups et cubes

Deux méthodes très puissantes pour agréger les données sont rollup() et cube(). Avec un groupBy() simple, vous n’obtenez que les agrégations existantes, tandis que rollup() et cube() créent des structures hiérarchiques, donc des niveaux supplémentaires d’agrégation.

Par exemple, rollup() agrège de gauche à droite, en affichant chaque permutation possible au fil des itérations. Pour chaque département, on voit chaque personne du service ainsi que le total du groupe.

De son côté, cube() affiche toutes les permutations possibles pour les colonnes agrégées : chaque département séparément, chaque employé séparément et toutes les combinaisons des deux. 

En résumé :

  • Rollup crée des sous-totaux hiérarchiques selon l’ordre des colonnes.
  • Cube génère des sous-totaux pour toutes les combinaisons possibles des colonnes spécifiées.

Exemple de code et de sortie ci-dessous :

rollup() Code :

# using rollup
df.rollup("department", “employee”).sum("salary").show()

rollup() Résultat :

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|    NULL|      23700|
|     Sales|   Alice|       5000|
|     Sales|    NULL|       9800|
|     Sales|     Bob|       4800|
|        HR|   Carol|       4000|
|        HR|    NULL|       7900|
|        HR|   David|       3900|
|        IT|     Eve|       6000|
|        IT|    NULL|       6000|
+----------+--------+-----------+

cube() Code :

#using cube
df.cube("department", "employee").sum("salary").show()

cube() Résultat :

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|   Alice|       5000|
|      NULL|    NULL|      23700|
|     Sales|   Alice|       5000|
|     Sales|    NULL|       9800|
|     Sales|     Bob|       4800|
|      NULL|     Bob|       4800|
|      NULL|   Carol|       4000|
|        HR|   Carol|       4000|
|        HR|    NULL|       7900|
|        HR|   David|       3900|
|      NULL|   David|       3900|
|        IT|     Eve|       6000|
|        IT|    NULL|       6000|
|      NULL|     Eve|       6000|
+----------+--------+-----------+

Grouping sets

Les grouping sets permettent de définir plusieurs niveaux d’agrégation. Dans cet exemple, j’agrège au niveau département+employé, département seul, et un total général pour le salaire.

Vous remarquerez une syntaxe un peu différente avec groupingSets(). On définit d’abord la liste des sets [(“department”, “employee”), (“department”, ), ()] : le premier correspond à département+employé, le second aux seuls départements, et le dernier () à « tous ». 

On définit ensuite les colonnes d’agrégation dans ces sets. Le reste se fait comme d’habitude avec .agg().

df.groupingSets([("department", "employee"), ("department",), ()], "department","employee").agg(sf.sum("salary")).sort("department","employee").show()

Résultat : 

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|    NULL|      23700|
|        HR|    NULL|       7900|
|        HR|   Carol|       4000|
|        HR|   David|       3900|
|        IT|    NULL|       6000|
|        IT|     Eve|       6000|
|     Sales|    NULL|       9800|
|     Sales|   Alice|       5000|
|     Sales|     Bob|       4800|
+----------+--------+-----------+

Fonctions d’agrégation personnalisées

Enfin, vous pouvez créer des fonctions d’agrégation personnalisées via les User Defined Functions (UDFs) en PySpark. Deux options existent : udf et pandas_udf. Toutes deux permettent de créer vos propres fonctions, chacune avec ses avantages et limites.

Le udf standard de Spark crée des fonctions natives Spark. Cela impose une syntaxe et des types de données compatibles Spark, ce qui peut restreindre certaines opérations, mais tire pleinement parti de la puissance de calcul distribuée de Spark — idéal pour les grands jeux de données.

from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType

@udf(returnType=IntegerType())
def bonus(salary):
    return int(salary * 0.1)

df.withColumn("bonus", bonus(df.salary)).show()

À l’inverse, pandas_udf permet de créer des fonctions personnalisées plus « pythoniques ». Vous êtes alors limité par ce que Pandas permet, et non plus uniquement Spark. 

En revanche, vous ne profitez pas de la distribution de calcul de Spark et vous vous appuyez sur des calculs locaux. C’est mieux adapté à des jeux de données petits à moyens.

import pandas as pd
from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def salary_bonus(s: pd.Series) -> pd.Series:
    return s * 0.1

df.withColumn("bonus", salary_bonus(df.salary)).show()

Globalement, n’utilisez des fonctions personnalisées que si aucune autre solution ne permet l’agrégation souhaitée. Inutile de réinventer la roue ! 

Par ailleurs, si une opération peut être vectorisée (comme ci-dessus), privilégiez une version vectorisée plutôt qu’un UDF d’agrégation. Pour aller plus loin sur groupBy(), lisez cet article plus détaillé sur le cadre split-apply-combine avec pandas.

Filtrer des données agrégées

Avec PySpark, vous pouvez filtrer des groupes sur la base de mesures agrégées après le regroupement via filter(). Vous fournissez alors une condition en Python ou en expressions SQL. 

À noter : vous pouvez aussi utiliser where(), alias de filter() réalisant la même opération.

Vous pouvez filtrer avant ou après l’agrégation. Le filtrage en amont limite les données agrégées et peut améliorer les performances. 

Par exemple, peut-être souhaitez-vous uniquement compter les employés au-dessus d’un certain salaire pour identifier les employés « haut salaire ».

# filtrer les salaires élevés
filter_df = df.filter(df.salary > 4000)

# agréger et compter les employés à haut salaire
agg_filter_df = filter_df.groupBy("department").agg(sf.count("*").alias("high_salary_emp"))
agg_filter_df.show()

En revanche, filtrer après l’agrégation n’affecte pas l’agrégation initiale. Vous pouvez compter tous les employés d’un département puis n’afficher que les départements au-dessus d’une certaine taille. Cela ajoute une étape, sans pour autant alourdir fortement la charge.

# créer un dataframe comptant le nombre d’employés
agg_df = df.groupBy("department").agg(sf.count("*").alias("num_employees"))

# filtrer les départements avec plus d’un employé
agg_df.where("num_employees > 1").show()

Choisissez le moment du filtrage avec discernement, car il impacte l’exactitude et le résultat final. Filtrer avant peut réduire le total final, tandis que filtrer après peut laisser trop d’éléments dans la réponse.

Optimiser les performances de groupBy en PySpark

Même si PySpark optimise automatiquement groupBy, certaines stratégies peuvent encore accélérer le traitement. Minimiser le shuffle, atténuer le skew et optimiser l’exécution contribuent tous à améliorer les performances de groupBy.

Gestion des shuffles

groupBy provoque un shuffle, qui redistribue les données entre partitions pour regrouper des clés similaires. Les shuffles sollicitent fortement le réseau et le disque, il est donc crucial de les minimiser ou de les optimiser. Techniques clés :

  • Utiliser repartition() à bon escient : indiquez à PySpark la colonne de partition pour réduire l’espace de recherche
  • Ajuster spark.sql.shuffle.partitions : par défaut à 200. Diminuez pour de petits jeux de données.
  • Activer la compression des shuffles pour réduire la charge réseau et disque.

Exemples d’améliorations de la gestion des shuffles :

# réduire le nombre de partitions de shuffle
spark.conf.set("spark.sql.shuffle.partitions", "64")  # à ajuster selon la taille du cluster

# s’assurer que Spark compresse les données
spark.conf.set("spark.shuffle.compress", "true") # compression des transferts réseau
spark.conf.set("spark.shuffle.spill.compress", "true") # compression des débordements disque

# repartitionner avant le groupement pour optimiser l’agrégation
df.repartition("department").groupBy("department").sum("salary").show()

Techniques de mitigation du skew

Le skew survient quand certaines clés sont surreprésentées, entraînant des charges de partitions inégales et des tâches lentes. Certains workers sont davantage sollicités. On peut atténuer le skew via le salting, la réduction des jointures biaisées et les broadcast joins pour les petites dimensions.

  • Salting : ajouter une colonne de nombres aléatoires pour répartir plus uniformément la charge
  • Repartitionner pour réduire le skew : forcer Spark à utiliser une autre colonne de partition peut mieux équilibrer la charge
# Exemple de salting :

df = df.withColumn("salted_key", sf.rand()) # Crée une colonne aléatoire
df = df.repartition(2, 'salted_key') # repartitionner via cette clé
df.groupBy(sf.spark_partition_id()).count().show()

Optimisation de l’exécution

Les plans logique et physique des jobs PySpark peuvent être optimisés via des fonctionnalités intégrées. Elles améliorent les performances de vos agrégations groupBy.

  • Catalyst optimizer : Spark réécrit automatiquement les plans inefficaces. Des transformations déclaratives (plutôt que des boucles procédurales) l’aident à mieux travailler.
  • Mise en cache : utile lorsque le même résultat de groupBy est réutilisé plusieurs fois dans un pipeline.
  • Broadcast joins : diffuser les petits DataFrames en mémoire tandis que les grands sont partitionnés pour minimiser les surcoûts réseau et disque.
grouped_df = df.groupBy("department").sum("salary").cache()
grouped_df.show() # vous pouvez mettre en cache des résultats intermédiaires

# utiliser le broadcast (pseudo-code) :
from pyspark.sql.functions import broadcast
df.join(broadcast(smaller_df), "department").show()

Pour plus de détails sur l’optimisation des opérations PySpark, lisez cet article sur les PySpark Joins afin de comprendre une partie des mécanismes internes.

Analyse comparative avec les opérations RDD

L’API DataFrame est recommandée pour la plupart des utilisateurs PySpark en raison de son niveau d’abstraction et de son efficacité (grâce au Catalyst Optimizer). Comprendre la couche RDD (Resilient Distributed Dataset) reste utile, notamment pour davantage de contrôle ou lors de migrations d’ancien code Spark. En général, l’API DataFrame ou SQL offre de meilleures performances.

Le RDD est le cœur de PySpark, sur lequel tout (y compris l’API DataFrame) est construit. Il travaille majoritairement en mémoire et gère parfois mieux des flux que l’API DataFrame. Les données sont toutefois immuables : on ne peut pas modifier un RDD après sa création. 

Cette section compare groupBy dans l’API DataFrame à ses équivalents RDD, en termes de performance et d’adéquation aux usages.

Commençons par convertir le DataFrame en RDD.

# Convertir le DataFrame en RDD de (clé, valeur)
rdd = df[['department','salary']].rdd
rdd.collect()

Nous pouvons maintenant utiliser groupByKey() et reduceByKey() sur le RDD. Commençons par groupByKey().

La méthode groupByKey() regroupe sur un même exécuteur toutes les valeurs partageant la même clé. Cela peut prendre du temps, car beaucoup de données doivent être déplacées, et cela peut accentuer le skew.

rdd.groupByKey().mapValues(list).collect()

À la place, utilisons reduceByKey(). Cette méthode agrège d’abord au sein de chaque partition, puis déplace les résultats sur le réseau. Elle minimise le volume de données shuffle et fonctionne donc plus vite. 

Ainsi, reduceByKey est fortement préférable à groupByKey car elle réduit significativement les opérations de shuffle.

Concrètement, on agrège dans chaque partition, on combine, puis on agrège globalement. C’est idéal pour des agrégations rassemblant les clés comme sum() et max().

from operator import add

rdd.reduceByKey(add).collect()

Bénéfices de performance : DataFrame vs RDD

Voici un tableau récapitulant les avantages de chaque méthode d’agrégation.

Caractéristique

DataFrame groupBy

RDD groupByKey()

RDD reduceByKey()

Niveau d’abstraction

Élevé

Faible

Faible

Exécution optimisée

Oui (Catalyst & Tungsten)

Non

Non

Minimisation du shuffle

Oui

❌ (Shuffle complet)

✅ (Avec combiner)

Efficacité mémoire

Élevée

Faible

Moyenne–élevée

Flexibilité

Modérée

Élevée

Moyenne

Performance sur big data

Excellente

Faible

Bonne

Cas recommandé

La plupart des cas

Rares/spécifiques

Agrégations bas niveau sur mesure

Le cours Big Data Fundamentals with PySpark couvre en détail la programmation avec les RDDs et explique leur rôle fondamental dans PySpark.

Requêtes PySpark SQL GROUP BY

Une autre manière d’agréger dans PySpark consiste à utiliser l’API SQL pour écrire des instructions en SQL. Idéal si vous êtes plus à l’aise avec ce langage.

Commencez par créer une vue temporaire via createOrReplaceTempView() sur le DataFrame. Vous pouvez ensuite écrire votre instruction avec spark.sql().

# Créer une vue temporaire à partir du DataFrame
df.createOrReplaceTempView("employees")

# Écrire une requête de type SQL
spark.sql("""
    SELECT department, AVG(salary) AS avg_salary
    FROM employees
    GROUP BY department
""").show()

Résultat :

+----------+----------+
|department|avg_salary|
+----------+----------+
|     Sales|    4900.0|
|        HR|    3950.0|
|        IT|    6000.0|
+----------+----------+

Comme vous le voyez, c’est aussi simple qu’une requête SQL classique. L’API SQL est pertinente pour des raisons de familiarité et de facilité d’écriture. L’API DataFrame est souvent préférable : elle assure des types cohérents et fournit en général de meilleures optimisations. 

Vous accédez à des méthodes DataFrame qui ne sont pas disponibles via Spark SQL. Pour en savoir plus sur GROUP BY en SQL, lisez cet article sur GROUP BY et HAVING en SQL.

Cas d’usage concrets

Les cas d’usage de l’agrégation sont nombreux. En pratique, vous aurez presque toujours besoin d’agréger vos données pour les analyser et les partager.

Voici quelques exemples. Certains sont du pseudo-code ou ne s’appliquent pas exactement à notre jeu de test, mais illustrent la manière d’écrire ces agrégations.

Business intelligence

Un thème récurrent est l’analyse hiérarchique avec groupBy. Imaginez une colonne « revenue » et l’envie d’évaluer la performance de chaque département. On groupe par département et on évalue sa contribution au revenu.

CopyEdit
df.groupBy("department").agg(sum("revenue").alias("total_revenue")).show()

Analyse de séries temporelles

Les agrégations donnent accès aux fonctions de fenêtre. Elles examinent une fenêtre glissante de données, soit une suite de lignes, soit une période temporelle. 

On définit d’abord une fenêtre avec l’objet Window pour préciser partitonBy et orderBy. On utilise ensuite cette fenêtre dans la méthode over() de la fonction d’agrégation choisie. Par exemple, une moyenne glissante ressemble à sf.avg().over(window).show().

Imaginons que l’on souhaite la moyenne glissante des salaires par département. On groupe par department et on ordonne par une nouvelle colonne employeeId. Les IDs étant souvent séquentiels, on observe l’évolution de la moyenne salariale dans le temps. 

from pyspark.sql.window import Window

# Définir la fenêtre (partition et ordre)
windowSpec = Window.partitionBy("department").orderBy("employeeId")

# Appliquer sf.avg sur "salary" sur la fenêtre définie
df.withColumn("rolling_avg", sf.avg("salary").over(windowSpec)).show()

C’est très proche des fonctions de fenêtre en SQL.

Analytique média

Supposons que nous voulions comprendre différents indicateurs utilisateur pour une entreprise média : temps de visionnage total et nombre de vidéos uniques. On peut groupBy() sur user_id puis utiliser .agg() avec sum() et countDistinct() pour différents métriques.

df.groupBy("user_id").agg(
    sf.sum("watch_time").alias("total_watch"), # somme des minutes regardées
    sf.countDistinct("video_id").alias("unique_views") # nombre de vidéos uniques
).show()

Bonnes pratiques et pièges à éviter avec PySpark GroupBy

En agrégeant des données dans PySpark, on tombe facilement dans des pièges menant à des performances médiocres et des temps d’exécution longs. Ces bonnes pratiques garantissent performance, justesse et scalabilité. 

Pièges fréquents

Voici des façons courantes de transformer une requête simple en programme interminable.

  1. Surutilisation de groupByKey() dans les RDD :
    • Évitez groupByKey() sur RDD : il force PySpark à faire du shuffle entre partitions, avec un fort coût réseau et disque.
    • Préférez reduceByKey() ou restez sur l’API DataFrame.
  2. Ne pas gérer le skew :
    • Si un groupe domine (ex. : un département avec des millions d’enregistrements), ce task devient un goulot d’étranglement.
    • Utilisez le salting ou un partitionnement personnalisé via repartition() pour mieux répartir la charge.
  3. Oublier de chaîner les fonctions d’agrégation :
    • Écrire des agrégations séparées (ex. : sum() dans un groupBy puis count() sur la même base groupée dans une autre ligne) fragmente l’optimisation de PySpark.
    • PySpark utilise une évaluation paresseuse : il génère le plan entier avant les actions. Chaînez les transformations (ex. : groupBy().agg().filter()) pour permettre à Catalyst d’optimiser.
  4. Mauvaise utilisation des UDF :
    • Les UDF Python empêchent Spark d’optimiser totalement les requêtes, Catalyst ne pouvant pas les optimiser.
    • Si possible, utilisez les fonctions intégrées de Spark ou des pandas UDF pour de meilleures performances.
  5. Mauvaise gestion mémoire :
    • Les agrégations sont coûteuses et gourmandes en mémoire. La réutilisation non maîtrisée des résultats peut provoquer des saturations.
    • Surveillez l’usage mémoire et envisagez persist() ou cache() lors de réutilisations.

Astuces d’optimisation et d’efficacité

Voici quelques leviers pour garder vos agrégations PySpark fluides !

Gestion des valeurs nulles

  • Agréger sur des colonnes contenant des nulls peut donner des résultats inattendus.
  • Utilisez na.fill() ou na.drop() avant l’agrégation.
df.na.fill({"salary": 0}).groupBy("department").sum("salary").show()

Éviter les shuffles excessifs

  • Comme expliqué plus haut, les shuffles ralentissent. Repartitionnez logiquement avant le groupement pour limiter leur taille.
  • Ajustez spark.sql.shuffle.partitions selon le volume de données pour trouver le bon nombre de partitions.

Utiliser les plans d’exécution

  • Surveillez les plans via .explain() pour comprendre l’exécution physique.
  • Repérez les shuffles larges, les hints de broadcast et les scans inefficaces
df.groupBy("department").sum("salary").explain(True)

Résultat :

	== Parsed Logical Plan ==
	'Aggregate ['department], ['department, unresolvedalias('sum(salary#600L))]
	+- LogicalRDD [department#598, employee#599, salary#600L], false

	== Analyzed Logical Plan ==
	department: string, sum(salary): bigint
	Aggregate [department#598], [department#598, sum(salary#600L) AS sum(salary)#679L]
	+- LogicalRDD [department#598, employee#599, salary#600L], false

	== Optimized Logical Plan ==
	Aggregate [department#598], [department#598, sum(salary#600L) AS sum(salary)#679L]
	+- Project [department#598, salary#600L]
	   +- LogicalRDD [department#598, employee#599, salary#600L], false

	== Physical Plan ==
	AdaptiveSparkPlan isFinalPlan=false
	+- HashAggregate(keys=[department#598], functions=[sum(salary#600L)], output=[department#598, sum(salary)#679L])
	   +- Exchange hashpartitioning(department#598, 200), ENSURE_REQUIREMENTS, [plan_id=1841]
	      +- HashAggregate(keys=[department#598], functions=[partial_sum(salary#600L)], output=[department#598, sum#681L])
	         +- Project [department#598, salary#600L]
	            +- Scan ExistingRDD[department#598,employee#599,salary#600L]
	

Surveiller les performances avec l’UI Spark

  • Utilisez l’interface Web Spark pour suivre les stages, les tâches et identifier des agrégations lentes ou du skew.
  • Apprenez-en plus sur l’UI Spark dans ce cours Introduction to Spark SQL in Python.

Exploiter le tuning de configuration

Spark propose de nombreux paramètres à ajuster pour des charges groupBy importantes. Voir le tableau ci-dessous pour quelques clés.

Configuration

Description

spark.sql.shuffle.partitions

Contrôle le nombre de partitions pour les shuffles. Réduisez pour de petits jobs, augmentez pour de gros volumes. (Par défaut : 200)

spark.sql.autoBroadcastJoinThreshold

Active la diffusion automatique des petites tables. Mettez à -1 pour désactiver ou augmentez pour supporter de plus grosses jointures. (Par défaut : 10 Mo)

spark.executor.memory

Mémoire disponible par exécuteur. Augmentez-la pour de grosses agrégations. (Par défaut : 4g)

spark.sql.adaptive.enabled

Active l’Adaptive Query Execution (AQE), qui optimise dynamiquement shuffles, skew et jointures. (Par défaut : true)

Notes et détails d’implémentation avancés

Pour les utilisateurs avancés ou en contexte de très grands volumes, comprendre le fonctionnement interne de groupBy est essentiel. Voici quelques points clés.

Évaluation paresseuse et plan d’exécution

Toutes les opérations DataFrame en PySpark, y compris groupBy, sont évaluées paresseusement. Aucune computation n’a lieu tant qu’une action (show(), collect(), etc.) n’est pas déclenchée. 

Catalyst peut ainsi réorganiser, combiner ou supprimer des opérations pour améliorer les performances. Profitez-en et chaînez plusieurs agrégations et méthodes pour permettre la création d’un plan optimisé.

Type de retour et collisions de noms

L’appel à groupBy() renvoie un objet GroupedData. Exécuter une fonction d’agrégation comme sum() renvoie ensuite un objet de type DataFrame. Avec show(), les résultats s’affichent.

N’oubliez pas que groupBy().agg() renvoie un nouveau DataFrame avec de nouvelles colonnes. Utilisez systématiquement des alias via alias() pour éviter les collisions de noms et clarifier les sorties. Sans alias, vous risquez des collisions de colonnes ou des problèmes de jointure plus tard.

from pyspark.sql import functions as F

df.groupBy("department") \

  .agg(

      sum("salary").alias("total_salary"),

      count("*").alias("employee_count")

  )

Format des résultats d’agrégation

Notez qu’en PySpark, l’ordre des lignes n’est pas garanti dans les résultats d’agrégation. Utilisez orderBy() pour des sorties prévisibles et un tri cohérent.

df.groupBy("department").sum("salary").orderBy("department").show()

Comportements et limites spécifiques aux versions

Certaines fonctionnalités et limitations varient selon la version de PySpark. Catalyst a été introduit en 1.3 puis fortement amélioré en 2.0. Des méthodes comme groupingSets() sont apparues dans des versions ultérieures. 

Voici un tableau récapitulant quelques évolutions majeures. Assurez-vous d’utiliser la version de Spark et PySpark offrant les fonctionnalités dont vous avez besoin ! 

Fonctionnalité

Version Spark

Description

cube() / rollup()

1.4+

Agrégations hiérarchiques utiles pour l’analytique OLAP (3.4+ prend en charge Spark Connect)

pandas_udf

2.3+

UDF vectorisées via Apache Arrow pour une exécution plus rapide (3.4+ prend en charge Spark Connect, 4.0+ prend en charge SCALAR)

Adaptive Query Execution (AQE)

3.0+

Ajuste dynamiquement les jointures, shuffles et la gestion du skew à l’exécution

Mode de compatibilité ANSI SQL

3.0+

Des erreurs et des expressions plus conformes et précises

groupingSets()

4.0+

Permet plusieurs groupements en une seule agrégation (ex. : sous-totaux)

Conclusion

La fonction groupBy de PySpark est un outil essentiel pour l’agrégation de données en environnement distribué. Qu’il s’agisse de synthétiser par région, de calculer des moyennes ou de conduire des analyses multi-niveaux complexes, groupBy offre une API scalable et flexible pour les charges big data.

Gardez à l’esprit ces bonnes pratiques pour réduire le temps d’exécution et les coûts :

  • Utilisez les fonctions intégrées dès que possible
  • Évitez les pièges de performance comme les shuffles excessifs et le skew.
  • Exploitez les optimisations et outils de profilage de Spark pour affiner les gros jobs.

À mesure que PySpark évolue, on peut s’attendre à des optimisations plus intelligentes, un support natif des agrégations complexes et une meilleure intégration avec Pandas et la syntaxe SQL. L’intégration avec des modèles plus larges pourrait aussi s’améliorer. Pour aller plus loin avec PySpark, explorez ces ressources DataCamp :

PySpark groupBy : FAQ

Que renvoie réellement la méthode groupBy() de PySpark ?

Elle renvoie un objet GroupedData, et non un DataFrame. Cet objet doit être suivi d’une méthode d’agrégation comme .count(), .sum() ou .agg() pour obtenir un nouveau DataFrame avec les résultats groupés.

Qu’est-ce que le data skew et quel est son impact sur groupBy() ?

Le skew des données survient lorsqu’une ou quelques clés concentrent une part disproportionnée des données, surchargeant certains exécutants pendant que d’autres restent inactifs. Cela peut dégrader les performances, voire faire échouer un job. Atténuez-le via le salting, un partitionnement personnalisé ou des broadcast joins pour les petites dimensions.

Puis-je utiliser groupBy() sur plusieurs colonnes dans PySpark ?

Oui, vous pouvez regrouper par une liste de colonnes.

Quand utiliser des fonctions d’agrégation personnalisées (UDF ou pandas_udf) ?

N’utilisez-les que lorsque les fonctions intégrées ne suffisent pas. Les fonctions natives sont plus rapides et bénéficient des optimisations de Spark. Les UDF empêchent Catalyst d’optimiser pleinement le plan, tandis que pandas_udf privilégie la flexibilité au détriment de la scalabilité et convient surtout à des jeux de données moyens.

Comment optimiser les opérations groupBy() pour la performance dans PySpark ?

Minimisez les shuffles avec .repartition() et en ajustant spark.sql.shuffle.partitions. Pensez à .cache() si vous réutilisez des résultats agrégés. Analysez l’optimisation de Spark et ses recommandations avec .explain().


Tim Lu's photo
Author
Tim Lu
LinkedIn

Je suis un data scientist avec de l'expérience dans l'analyse spatiale, l'apprentissage automatique et les pipelines de données. J'ai travaillé avec GCP, Hadoop, Hive, Snowflake, Airflow et d'autres processus d'ingénierie et de science des données.

Sujets
PySpark
Python

Meilleurs cours PySpark

Cours

Principes fondamentaux de PySpark

4 h
157.8K
Apprenez à mettre en œuvre la gestion des données distribuées et l'apprentissage automatique dans Spark à l'aide du package PySpark.
Afficher les détailsRight Arrow
Commencer Le Cours
Voir plusRight Arrow