Aller au contenu principal
big-dataDécouper le calcul

Découper le calcul

Ce que ce chapitre apporte

  • Écrire un calcul sous forme map puis reduce, et savoir ce que chaque étape produit.
  • Distinguer une opération qui se recolle d'une opération qui ne se recolle pas.
  • Transporter l'état suffisant pour reconstituer une moyenne, et savoir pourquoi il en faut deux morceaux.
  • Mesurer et réduire le volume déplacé lors du mélange.
  • Reconnaître un déséquilibre de charge et remonter à la clé qui le provoque.

Une seule machine ne suffit plus : on coupe les données, chaque machine calcule sur son morceau, et l'on rassemble. Dit ainsi, c'est une évidence. En pratique, deux choses seulement rendent l'exercice difficile, et ce sont les deux seules qui méritent d'être apprises. La moyenne de deux moyennes n'est pas la moyenne. Et les morceaux ne sont jamais de la même taille.

Le chapitre précédent a montré comment traiter des données qui ne tiennent pas en mémoire, sur une seule machine, en une seule lecture. Cette solution atteint sa limite quand la lecture elle-même devient trop longue : lire un pétaoctet à 200 mégaoctets par seconde demande deux mois, et aucun algorithme en une passe n'y changera rien.

La seule issue est alors de lire à plusieurs endroits en même temps. C'est ce que fait le motif présenté ici, qui porte le nom de ses deux étapes, map et reduce, et qui structure la quasi-totalité des systèmes de traitement réparti, quel que soit leur habillage.

Le motif, en quatre temps

Map, mélange, reduce

Map transforme chaque élément, indépendamment des autres, en un ou plusieurs couples clé-valeur. C'est l'étape qui se parallélise sans effort, puisque aucune machine n'a besoin de savoir ce que font les autres.

Le mélange regroupe tous les couples qui portent la même clé sur une même machine. C'est la seule étape où les données traversent le réseau, et c'est donc elle qui coûte.

Reduce applique une opération à toutes les valeurs d'une même clé, et produit un résultat par clé.

Le comptage de mots est l'exemple canonique du domaine, et il le mérite : il tient en trois lignes et contient tout le motif.

somme · par tourniquet

Les données sont coupées en morceaux, un par machine.

machine 1

  • le chat dort sur le mur

machine 2

  • le chien dort sous le banc

machine 3

  • le chat court après le chien
Compter les mots d'un texte sur trois machines

Les quatre phases se parcourent avec les boutons. La découpe répartit les lignes, sans regarder ce qu'elles contiennent. Le map produit un couple par mot, (mot, 1), sur chaque machine indépendamment. Le mélange rassemble par mot, et l'on voit alors les couleurs se mêler : ce qui vient de la machine 1 rejoint ce qui vient de la machine 3. La réduction additionne.

Le map ne sait rien, et c'est ce qui le rend parallélisable

Le traitement d'une ligne par le map ne dépend d'aucune autre ligne. C'est cette propriété, et elle seule, qui autorise à découper n'importe comment et à traiter en parallèle sans coordination.

Toute logique qui exige de connaître les autres lignes doit donc être repoussée dans le reduce, après le regroupement par clé. Un map qui a besoin d'un compteur global n'est pas un map.

Écrit à la main, le motif ne fait qu'une vingtaine de lignes de Python. Rien n'y est mystérieux, et le voir en entier vaut mieux que de faire confiance à un moteur.

main.py
Sortie
>_ Prêt à exécuter…

Le mélange est ce qui coûte

Dans la figure ci-dessus, dix-huit couples traversent le réseau pour produire dix résultats. Sur un texte réel, c'est un couple par mot du corpus : des milliards de couples de la forme (le, 1), tous identiques, tous transportés séparément.

L'idée qui corrige cela est simple : réduire d'abord localement, sur chaque machine, avant d'envoyer.

Combineur

Un combineur est une réduction appliquée localement, sur chaque morceau, avant le mélange. Il ne change pas le résultat, il change ce qui circule.

Une machine qui a vu deux cents fois le mot le envoie (le, 200) au lieu de deux cents couples (le, 1). Le résultat final est identique, le volume transporté est divisé par deux cents.

somme · par tourniquet · combineur local

Les données sont coupées en morceaux, un par machine.

machine 1

  • le chat dort sur le mur

machine 2

  • le chien dort sous le banc

machine 3

  • le chat court après le chien
Le même comptage, avec réduction locale avant l'envoi

La phase de mélange annonce maintenant dix couples au lieu de dix-huit. Sur trois phrases le gain est anecdotique ; sur un corpus, il fait la différence entre un traitement qui tient et un traitement qui sature le réseau.

Un combineur n'est pas une optimisation gratuite

Il n'est légitime que si la réduction admet d'être appliquée en deux temps sans changer le résultat. Pour une somme ou un maximum, c'est vrai. Pour une moyenne, c'est faux, et la section suivante montre à quel point.

Un moteur de traitement réparti applique le combineur quand il le juge utile, sans prévenir, et parfois pas du tout. Un calcul dont le résultat dépend de l'application ou non du combineur est donc un calcul faux, dont le résultat varie d'une exécution à l'autre sans que rien ne le signale.

Ce qui se recolle, et ce qui ne se recolle pas

Opération associative

Une opération est associative quand l'ordre des regroupements ne change pas le résultat : (a ⊕ b) ⊕ c vaut a ⊕ (b ⊕ c).

C'est exactement la condition pour qu'une réduction se calcule par morceaux. La somme, le produit, le minimum, le maximum et le comptage sont associatifs. Une somme de sommes partielles est une somme ; un maximum de maximums partiels est un maximum.

La moyenne ne l'est pas.

L'affirmation se vérifie sur quatre valeurs.

moyenne · par tourniquet · combineur naïf

Les données sont coupées en morceaux, un par machine.

machine 1

  • maths 12
  • maths 18

machine 2

  • maths 12
  • info 10
Moyenne par matière, avec un combineur écrit trop vite

La machine 1 a reçu 12 et 18, et envoie leur moyenne, 15. La machine 2 a reçu 12, et envoie 12. Le réducteur fait la moyenne de 15 et de 12, soit 13,5. La bonne réponse est 14, et la figure l'affiche à côté.

L'erreur n'est pas dans l'arithmétique : elle est dans ce qui a été transporté. Une moyenne résume deux informations, un total et un effectif, et n'en envoyer qu'une seule revient à traiter un morceau de deux valeurs comme un morceau d'une seule.

Transporter de quoi refaire le calcul, et ne diviser qu'à la fin

Le combineur juste envoie le couple (somme, effectif). Le réducteur additionne les sommes, additionne les effectifs, et divise une seule fois, à la fin.

Ce raisonnement se généralise : une réduction non associative devient calculable par morceaux dès qu'on transporte un état intermédiaire qui, lui, se combine associativement.

Pour la variance, cet état comporte trois nombres, ceux du chapitre précédent. Pour un écart type, les mêmes. Pour une médiane, aucun état de taille fixe ne convient, ce qui est une autre façon de dire qu'elle n'est pas décomposable.

moyenne · par tourniquet · combineur local

Les données sont coupées en morceaux, un par machine.

machine 1

  • maths 12
  • maths 18

machine 2

  • maths 12
  • info 10
La même moyenne, avec un combineur qui transporte l'effectif
L'erreur ne se voit pas toujours, et c'est ce qui la rend dangereuse

Si les deux morceaux contiennent le même nombre de valeurs, la moyenne des moyennes donne la bonne réponse. Le calcul faux passe donc les essais sur des données bien réparties, et ne se met à mentir qu'en production, quand la répartition devient inégale.

Pire, l'écart est petit et vraisemblable : 13,5 au lieu de 14 ne choque personne et ne déclenche aucune alerte. C'est la signature des erreurs les plus coûteuses, celles qu'aucun symptôme ne trahit.

La parade tient en une question à poser sur chaque agrégat : puis-je le calculer sur deux moitiés et recoller le résultat ? Si la réponse est non, il faut transporter davantage.

Le même piège, écrit en Python, sur des données dont le déséquilibre est explicite.

main.py
Sortie
>_ Prêt à exécuter…

Les morceaux ne sont pas égaux

Le mélange impose que toutes les valeurs d'une clé finissent au même endroit. Répartir directement par clé évite donc de tout transporter deux fois, et c'est ce que font les moteurs réels. Cette économie a un prix.

somme · par clé

Les données sont coupées en morceaux, un par machine.

machine 1

  • paris 120
  • paris 80
  • paris 200
  • paris 60
  • paris 150
  • paris 90
  • lyon 40

machine 2

  • nice 30

machine 3

rien à traiter

Le morceau le plus chargé porte 2,63 fois sa part. Les autres machines l'attendront, et le calcul durera ce que dure la plus lente : découper ne sert à rien si les morceaux sont inégaux.

Ventes par ville, réparties selon la clé

Une machine porte presque tout, une autre presque rien, et la troisième reste vide. Le calcul durera ce que dure la plus chargée : ajouter des machines n'y change rien, puisque la clé paris ne se coupe pas en deux.

Déséquilibre de charge

Il y a déséquilibre quand la répartition par clé concentre une part disproportionnée des données sur une minorité de machines.

Le temps total d'un traitement réparti est celui de la machine la plus lente. Un déséquilibre de facteur dix annule donc le bénéfice de dix machines, et l'ajout de machines supplémentaires n'apporte plus rien.

Les données réelles sont presque toujours déséquilibrées

La répartition uniforme est l'exception, pas la règle. Les ventes se concentrent sur quelques produits, le trafic sur quelques pages, les connexions sur quelques heures, les messages sur quelques comptes. Une clé naturelle produit donc presque toujours un déséquilibre.

Trois parades, dans l'ordre de préférence :

Choisir une autre clé, plus fine, quand le calcul le permet. Compter par (ville, jour) plutôt que par ville multiplie le nombre de clés et disperse la charge.

Décomposer la clé dominante en ajoutant un suffixe tiré au hasard, réduire, puis réduire une seconde fois sur la clé nettoyée. Cela ne fonctionne que si la réduction est associative, ce qui ramène à la section précédente.

Traiter la clé dominante à part, avec un calcul spécifique. C'est la solution la moins élégante et souvent la plus rapide à mettre en place.

Où sont Hadoop et Spark

Le motif présenté ici est antérieur aux outils qui l'ont popularisé, et leur survivra. Il est utile de savoir ce que ces noms recouvrent, à condition de ne pas confondre le motif avec ses mises en œuvre.

Hadoop désigne un ensemble, dont un système de fichiers réparti qui découpe les fichiers en blocs répliqués sur plusieurs machines, et un moteur qui exécute des tâches map et reduce sur ces blocs, en écrivant sur disque entre chaque étape. Cette écriture systématique le rend robuste et lent.

Spark exécute le même motif en gardant les résultats intermédiaires en mémoire quand c'est possible, ce qui change les ordres de grandeur sur les traitements enchaînés. Son interface parle de transformations et d'actions plutôt que de map et de reduce, mais la mécanique du dessous est celle de ce chapitre, mélange compris.

Ce que ces outils ne dispensent pas de savoir

Aucun moteur ne détecte qu'une moyenne a été recollée à tort : il exécute ce qu'on lui décrit et rend un résultat vraisemblable.

Aucun moteur ne corrige un déséquilibre de clé : il peut le signaler dans ses journaux, et c'est tout.

Aucun moteur ne devine que le calcul demandé aurait pu tenir sur une seule machine en une passe. Un traitement réparti sur des données de quelques gigaoctets est souvent plus lent que la boucle du chapitre précédent, à cause du temps de mise en route et du mélange.

Exercices type

Exercice 1 : exprimer sous forme map et reduce le calcul du chiffre d'affaires total par magasin, à partir de lignes (magasin, produit, montant).

Afficher la solution

Le map émet (magasin, montant) pour chaque ligne, en ignorant le produit. Le reduce somme les montants d'une même clé.

La somme étant associative, un combineur local est légitime et divise le volume transporté par le nombre moyen de lignes par magasin et par morceau.

La clé est le magasin : s'il y en a quelques centaines et que le trafic est comparable, la répartition sera acceptable. S'il existe un magasin en ligne qui pèse la moitié du chiffre d'affaires, le déséquilibre est à prévoir.

Exercice 2 : même énoncé, mais on veut le panier moyen par magasin. Que doit transporter le combineur ?

Afficher la solution

Le couple (somme des montants, nombre de lignes).

Le reduce additionne les deux composantes séparément, puis divise une fois. Envoyer directement la moyenne locale donnerait la moyenne des moyennes, fausse dès que les morceaux n'ont pas le même nombre de lignes, c'est-à-dire presque toujours.

Le test à retenir : un état intermédiaire est valable si l'on peut le combiner associativement. (somme, effectif) se combine composante par composante ; une moyenne seule, non.

Exercice 3 : un traitement doit donner, par magasin, le montant médian. Que faire ?

Afficher la solution

La médiane n'est pas décomposable : aucun état de taille fixe ne permet de la recoller, et le motif ne s'applique donc pas directement.

Trois issues. Si les données d'un magasin tiennent en mémoire sur une machine, répartir par magasin et calculer la médiane exacte dans le reduce, qui voit alors toutes les valeurs de sa clé. C'est la solution la plus simple, et elle échoue sur le magasin dominant.

Sinon, transporter un histogramme des montants par tranches, qui se combine associativement en additionnant les effectifs des tranches, et lire la médiane dans l'histogramme cumulé. Le résultat est approché, à la largeur de tranche près, et cette erreur est connue et bornée.

Sinon, employer un algorithme de quantiles approchés, conçu pour cela.

Exercice 4 : un traitement sur 40 machines dure 50 minutes. Les journaux montrent que 39 machines ont fini en 3 minutes. Que conclure, et que tenter ?

Afficher la solution

Une clé concentre l'essentiel des données : c'est un déséquilibre franc, et non un problème de puissance.

Ajouter des machines ne changera rien, puisque la clé lourde reste indivisible. Le temps resterait de 50 minutes sur 400 machines.

Les pistes utiles : identifier la clé fautive en comptant les lignes par clé, ce qui est lui-même un map et reduce ; puis choisir une clé plus fine, décomposer la clé lourde par un suffixe aléatoire suivi d'une seconde réduction, ou traiter cette clé séparément.

Il arrive que la clé fautive soit une valeur par défaut : un inconnu, un null converti en chaîne, un identifiant à zéro. Dans ce cas, le déséquilibre est un symptôme de la qualité des données et non de leur distribution réelle.

Vérification rapideon peut se reprendre

1.Pourquoi l'étape map se parallélise-t-elle sans coordination ?

2.Quelle étape fait circuler les données sur le réseau ?

3.Un combineur qui envoie la moyenne locale de chaque morceau…

4.Quelle propriété rend une réduction calculable par morceaux ?

5.Un traitement sur 40 machines dure 50 minutes, dont 39 machines finissent en 3 minutes. La cause la plus probable ?

6.Que ne fait aucun moteur de traitement réparti à la place du concepteur ?

La méthode

  1. Écrire le map en vérifiant qu'il ne dépend d'aucune autre ligne.
  2. Nommer la clé, et se demander tout de suite combien de valeurs distinctes elle prend.
  3. Tester l'associativité de la réduction : puis-je calculer sur deux moitiés et recoller ?
  4. Si la réponse est non, chercher l'état intermédiaire qui, lui, se combine, et ne finir le calcul qu'à la fin.
  5. Activer un combineur dès que la réduction est associative, et mesurer ce qu'il économise.
  6. Regarder la répartition des clés avant de lancer, en comptant les lignes par clé.
  7. Lire les temps par machine après coup : un seul temps qui dépasse signale une clé, pas une panne.
  8. Se demander si le découpage était nécessaire : sous quelques dizaines de gigaoctets, une passe sur une machine gagne souvent.

Synthèse

  • Trois étapes : map indépendant, mélange par clé, reduce par clé.
  • Le map se parallélise parce qu'il ne sait rien des autres éléments.
  • Le mélange est la seule étape où les données traversent le réseau, donc la seule qui coûte vraiment.
  • Un combineur réduit localement avant d'envoyer et divise le volume transporté, sans changer le résultat.
  • Il n'est légitime que si la réduction est associative : somme, produit, minimum, maximum, comptage.
  • La moyenne n'est pas associative. La recoller à partir de moyennes partielles donne un résultat faux.
  • La parade consiste à transporter (somme, effectif) et à ne diviser qu'à la fin.
  • L'erreur disparaît quand les morceaux sont égaux, ce qui la fait passer les essais et échouer en production.
  • Répartir par clé évite un transport, et provoque un déséquilibre dès qu'une clé domine.
  • Le temps total est celui de la machine la plus lente : ajouter des machines ne divise pas une clé.
  • Une clé fautive est parfois une valeur par défaut, symptôme de qualité des données.
  • Hadoop et Spark exécutent ce motif ; aucun des deux ne détecte un agrégat mal recollé.