Pourquoi préférer reduceByKey à groupByKey ?

Questions d’entrevue Apache Spark

Intermédiaireshufflereducebykeyperformancerdd

La réponse courte

Parce que reduceByKey agrège avant le shuffle, et groupByKey transfère tout.

groupByKey déplace chaque valeur sur la machine qui détient la clé, puis vous agrégez. reduceByKey applique d’abord la fonction de réduction localement, sur chaque partition — l’équivalent d’un combiner MapReduce — et ne transfère que les résultats partiels.

Sur un million d’occurrences du prénom « Ali » réparties sur 200 partitions :

  • groupByKey transfère un million de valeurs.
  • reduceByKey transfère 200 sommes partielles.

C’est trois ordres de grandeur sur le volume réseau, pour un résultat identique.

Ce que l’intervieweur vérifie

Deux choses au-delà de la performance.

Que vous connaissez le risque de plantage. Avec groupByKey, toutes les valeurs d’une clé doivent tenir en mémoire sur un seul exécuteur. Une clé très fréquente provoque un OutOfMemoryError que reduceByKey n’aurait jamais produit. Ce n’est donc pas seulement plus lent : c’est plus fragile.

Que vous connaissez la condition d’usage. reduceByKey suppose une opération associative et commutative — somme, max, min, comptage. Elle impose aussi que le résultat ait le même type que les valeurs d’entrée, ce qui rend une moyenne impossible d’un seul coup.

Quand groupByKey reste légitime

Quand vous avez réellement besoin de la liste complète des valeurs : une concaténation ordonnée, un calcul de médiane, un traitement qui inspecte l’ensemble du groupe. Répondre « jamais » est une erreur — l’intervieweur attend la nuance.

Et pour le cas de la moyenne, la bonne réponse n’est ni l’un ni l’autre : c’est aggregateByKey, qui permet un type d’accumulateur différent du type des valeurs.

scala
val moyennes = notes
  .aggregateByKey((0.0, 0))(
    (acc, note) => (acc._1 + note, acc._2 + 1),
    (a, b) => (a._1 + b._1, a._2 + b._2),
  )
  .mapValues { case (somme, n) => somme / n }

La relance probable

« Et en DataFrame ? »

La question ne se pose plus : df.groupBy("cle").agg(sum("valeur")) passe par Catalyst, qui insère l’agrégation partielle tout seul. Le piège groupByKey est propre à l’API RDD — une bonne occasion de rappeler pourquoi le DataFrame est le choix par défaut.

Le détail, avec les quatre opérations voisines, est dans groupBy, groupByKey, reduceByKey et sortByKey.

Toutes les questions Apache Spark