MapReduce
1
Présentation (1/4)
• Pour exécuter un problème large de manière distribuée, il faut pouvoir
découper le problème en plusieurs problèmes de taille réduite à exécuter sur
chaque machine du cluster (stratégie algorithmique dite du divide and conquer
/ diviser pour régner).
• De multiples approches existent et ont existé pour cette division d'un problème
en plusieurs « sous-tâches ».
• MapReduce est un paradigme (un modèle) visant à généraliser les approches
existantes pour produire une approche unique applicable à tous les problèmes.
• MapReduce existait déjà depuis longtemps, notamment dans les langages
fonctionnels (Lisp, Scheme), mais la présentation du paradigme sous une forme
« rigoureuse », généralisable à tous les problèmes et orientée calcul distribué
est attribuable à un whitepaper issu du département de recherche de Google
publié en 2004 (« MapReduce: Simplified Data Processing on Large Clusters »).
2
Présentation (2/4)
MapReduce définit deux opérations distinctes à effectuer sur les données
d'entrée:
• La première, MAP, va transformer les données d'entrée en une série de couples
clef/valeur. Elle va regrouper les données en les associant à des clefs, choisies
de telle sorte que les couples clef/valeur aient un sens par rapport au problème à
résoudre. Par ailleurs, cette opération doit être parallélisable: on doit pouvoir
découper les données d'entrée en plusieurs fragments, et faire exécuter
l'opération MAP à chaque machine du cluster sur un fragment distinct.
• La seconde, REDUCE, va appliquer un traitement à toutes les valeurs de
chacune des clefs distinctes produite par l'opération MAP. Au terme de
l'opération REDUCE, on aura un résultat pour chacune des clefs distinctes.
Ici, on attribuera à chacune des machines du cluster une des clefs uniques
produites par MAP, en lui donnant la liste des valeurs associées à la clef.
Chacune des machines effectuera alors l'opération REDUCE pour cette clef.
3
Présentation (3/4)
On distingue donc 4 étapes distinctes dans un traitement MapReduce:
• Découper (split) les données d'entrée en plusieurs fragments.
• Mapper chacun de ces fragments pour obtenir des couples (clef ; valeur).
• Grouper (shuffle) ces couples (clef ; valeur) par clef.
• Réduire (reduce) les groupes indexés par clef en une forme finale, avec une
valeur pour chacune des clefs distinctes.
En modélisant le problème à résoudre de la sorte, on le rend parallélisable –
chacune de ces tâches à l'exception de la première seront effectuées de manière
distribuée.
4
Présentation (4/4)
Pour résoudre un problème via la méthodologie MapReduce avec Hadoop, on
devra donc:
• Choisir une manière de découper les données d'entrée de telle sorte que
l'opération MAP soit parallélisable.
• Définir quelle CLEF utiliser pour notre problème.
• Écrire le programme pour l'opération MAP.
• Écrire le programme pour l'opération REDUCE..
… et Hadoop se chargera du reste (problématiques calcul distribué,
groupement par clef distincte entre MAP et REDUCE, etc.).
5
Exemple concret (1/8)
• Imaginons qu'on nous donne un texte écrit en langue Française. On
souhaite déterminer pour un travail de recherche quels sont les mots les
plus utilisés au sein de ce texte (exemple Hadoop très répandu).
• Ici, nos données d'entrée sont constituées du contenu du texte.
• Première étape: déterminer une manière de découper (split) les données
d'entrée pour que chacune des machines puisse travailler sur une partie
du texte.
• Notre problème est ici très simple – on peut par exemple décider de
découper les données d'entrée ligne par ligne. Chacune des lignes du
texte sera un fragment de nos données d'entrée.
6
Exemple concret (2/8)
• Nos données d'entrée (le texte):
Celui qui croyait au ciel
Celui qui n'y croyait pas (Louis Aragon, La rose et le
[…] Réséda, 1943, fragment)
Fou qui fait le délicat
Fou qui songe à ses querelles
• Pour simplifier les choses, on va avant le découpage supprimer toute
ponctuation et tous les caractères accentués. On va également passer
l'intégralité du texte en minuscules.
7
Exemple concret (3/8)
• Nos données d'entrée (le texte):
celui qui croyait au ciel
celui qui ny croyait pas
fou qui fait le delicat
fou qui songe a ses querelles
• … on obtient 4 fragments depuis nos données d'entrée.
8
Exemple concret (4/8)
• On doit désormais déterminer la clef à utiliser pour notre opération
MAP, et écrire le code de l'opération MAP elle-même.
• Puisqu'on s'intéresse aux occurrences des mots dans le texte, et qu'à
terme on aura après l'opération REDUCE un résultat pour chacune des
clefs distinctes, la clef qui s'impose logiquement dans notre cas est: le
mot-lui même.
• Quand à notre opération MAP, elle sera elle aussi très simple: on va
simplement parcourir le fragment qui nous est fourni et, pour chacun
des mots, générer le couple clef/valeur: (MOT ; 1). La valeur indique ici
l’occurrence pour cette clef - puisqu'on a croisé le mot une fois, on
donne la valeur « 1 ».
9
Exemple concret (5/8)
• Le code de notre opération MAP sera donc (ici en pseudo code):
POUR MOT dans LIGNE, FAIRE :
GENERER COUPLE (MOT; 1)
• Pour chacun de nos fragments, les couples (clef; valeur) générés seront
donc:
celui qui croyait au ciel (celui;1) (qui;1) (croyait;1) (au;1) (ciel;1)
celui qui ny croyait pas (celui;1) (qui;1) (ny;1) (croyait;1) (pas;1)
fou qui fait le delicat (fou;1) (qui;1) (fait;1) (le;1) (delicat;1)
(fou;1) (qui;1) (songe;1) (a;1) (ses;1)
fou qui songe a ses querelles
(querelles;1)
28
Exemple concret (6/8)
• Une fois notre opération MAP effectuée (de manière distribuée),
Hadoop groupera (shuffle) tous les couples par clef commune.
• Cette opération est effectuée automatiquement par Hadoop. Elle est, là
aussi, effectuée de manière distribuée en utilisant un algorithme de tri
distribué, de manière récursive. Après son exécution, on obtiendra les 15
groupes suivants:
(celui;1) (celui;1) (fou;1) (fou;1) (fait;1) (le;1)
(qui;1) (qui;1) (qui;1) (qui;1) (delicat;1) (songe;1)
(croyait;1) (croyait;1) (a;1) (ses;1)
(au;1) (ciel;1) (ny;1) (pas;1) (querelles;1)
11
Exemple concret (7/8)
• Il nous reste à créer notre opération REDUCE, qui sera appelée pour
chacun des groupes/clef distincte.
• Dans notre cas, elle va simplement consister à additionner toutes les
valeurs liées à la clef spécifiée:
TOTAL=0
POUR COUPLE dans GROUPE, FAIRE:
TOTAL=TOTAL+1
RENVOYER TOTAL
12
Exemple concret (8/8)
• Une fois l'opération REDUCE effectuée, on obtiendra donc une valeur
unique pour chaque clef distincte. En l’occurrence, notre résultat sera:
qui : 4 •On constate que le mot le plus utilisé dans
celui : 2 notre texte est « qui », avec 4 occurrences,
croyait : 2 suivi de « celui », « croyait » et « fou », avec
fou : 2 2 occurrences chacun.
au : 1
ciel : 1
ny : 1
pas : 1
fait : 1
[…]
13
Exemple concret - Conclusion
• Notre exemple est évidemment trivial, et son exécution aurait été instantanée
même sur une machine unique, mais il est d'ores et déjà utile: on pourrait
tout à fait utiliser les mêmes implémentations de MAP et REDUCE sur
l'intégralité des textes d'une bibliothèque Française, et obtenir ainsi un bon
échantillon des mots les plus utilisés dans la langue Française.
• L’intérêt du modèle MapReduce est qu'il nous suffit de développer les deux
opérations réellement importantes du traitement: MAP et REDUCE, et de
bénéficier automatiquement de la possibilité d'effectuer le traitement sur un
nombre variable de machines de manière distribuée.
14
Schéma général
15
Posté par Mirko Krivanek : « What Is MapReduce? », credit @Tgrall
[Link]
34
Qu’est-ce que YARN ?
• YARN (Yet Another Resource Negociator) est un mécanisme dans Hadoop
permettant de gérer des travaux (jobs) sur un cluster de machines.
• YARN permet aux utilisateurs de lancer des jobs MapReduce sur des
données présentes dans HDFS, et de suivre (monitor) leur avancement,
récupérer les messages (logs) affichés par les programmes.
• Éventuellement YARN peut déplacer un processus d’une machine à l’autre
en cas de défaillance ou d’avancement jugé trop lent.
• En fait, YARN est transparent pour l’utilisateur. On lance l’exécution d’un
programme MapReduce et YARN fait en sorte qu’il soit exécuté le plus
rapidement possible.
YARN : MapReduce 2
• YARN est aussi appelé MRv2 (MapReduce 2). Ce n’est pas une refonte
mais une évolution du framework MapReduce.
• YARN répond aux problématiques suivantes du Map Reduce :
– Problème de limite de “Scalability” notamment par une meilleure
séparation de la gestion de l’état du cluster et des ressources.
~ 4000 Noeuds, 40 000 Tâches concourantes.
– Problème d’allocation des ressources.
18