Chapitre III :
MapReduce
PR NAOUFEL KRAIEM
Plan
MapReduce
Design Patterns MapReduce
2
MapReduce
Imaginons le probl me suivant :
Contexte : une cha ne de magasins dispers s travers le monde
Objectif : calculer le total des ventes par magasin
Supposons que toutes les ventes sont stock es dans un grand
livre sous la forme suivante :
Date
Ville
2016-01-01
London
2016-01-01
2016-01-02
&
Miami
Miami
&
Produit
Clothes
Music
Clothes
&
Prix
25.99
12.15
50.00
&
3
MapReduce
Solutions traditionnelles :
Pour chaque entr e, saisir la ville et le prix de vente. Si on trouve une entr e
avec une ville d j saisie, on les regroupe en faisant la somme des ventes.
M thode trop lente et peu efficace!
Dans un environnement de calcul traditionnel, on utilise g n ralement
des Hashtables, sous forme de : Clef, Valeur
Dans ce cas, la clef serait la ville, et la valeur le total des ventes
Date
Ville
Produit
2016-01-01
London
Clothes
2016-01-01
Miami
Music
2016-01-02
New York Toys
2016-01-02
Miami
Clothes
&
&
&
Prix
25.99
12.15
3,10
50.00
&
Cl
London
Miami
Valeur
25,99
62,15
New York
3,10
4
MapReduce
Si on utilise les hashtables sur 1To ?
Taille de la table => Probl me de m moire
Traitement s quentiel => temps de traitement trop long
a peut marcher quand m me!
Autres solutions : MapReduce !
Moyen plus efficace et plus rapide pour traiter les donn es
Au lieu davoir une seule personne qui parcourt le livre,
on recrute plusieurs ?
5
MapReduce : 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 ).
COURS BIG DATA - 2019
6
MapReduce : 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
faire ex cuter
les donn es d'entr e en plusieurs fragments, et
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.
COURS BIG DATA - 2019
7
MapReduce : 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.
COURS BIG DATA - 2019
8
MapReduce : 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.).
COURS BIG DATA - 2019
9
MapReduce : 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.
COURS BIG DATA - 2019
10
MapReduce : Exemple concret (2/8)
" Nos donn es d'entr e (le texte):
Celui qui croyait au ciel
Celui qui n'y croyait pas
[&]
Fou qui fait le d licat
Fou qui songe ses querelles
(Louis Aragon, La rose et le
R s da, 1943, fragment)
" 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.
COURS BIG DATA - 2019
11
Advertisement
MapReduce : 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.
COURS BIG DATA - 2019
12
MapReduce : 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
loccurrence pour cette clef - puisqu'on a crois le mot une fois, on
donne la valeur 1 .
COURS BIG DATA - 2019
13
MapReduce : 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 qui songe a ses querelles
(fou;1) (qui;1) (songe;1) (a;1) (ses;1)
(querelles;1)
28
Cours Big Data - 2019
MapReduce : 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)
Cours Big Data - 2019
15
MapReduce : 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
Cours Big Data - 2019
16
MapReduce : Exemple concret (8/8)
" Une fois l'op ration REDUCE effectu e, on obtiendra donc une valeur
unique pour chaque clef distincte. En loccurrence, notre r sultat sera:
qui : 4
celui : 2
croyait : 2
fou : 2
au : 1
ciel : 1
ny : 1
pas : 1
fait : 1
[&]
"On constate que le mot le plus utilis dans
notre texte est qui , avec 4 occurrences,
suivi de celui , croyait et fou , avec
2 occurrences chacun.
Cours Big Data - 2019
17
MapReduce : 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.
" Lint r t du mod le MapReduce est qu'il nous suffit de
d velopper les deux op rations r ellement importantes du
et de b n ficier
traitement: MAP et
automatiquement de la possibilit d'effectuer le traitement
sur un nombre variable de machines de mani re distribu e.
REDUCE,
Cours Big Data - 2019
18
MapReduce : Sch ma g n ral
19
MapReduce : Sch ma g n ral
Cours Big Data - 2019
20
MapReduce : Sch ma g n ral
Cours Big Data - 2019
21
Post par Mirko Krivanek : What Is MapReduce? , credit @Tgrall
http://www.datasciencecentral.com/forum/topics/what-is-map-reduce
34
MapReduce : Exemple Statistiques web
" Un autre exemple: on souhaite compter le nombre de visiteurs sur
chacune des pages d'un site Internet. On dispose des fichiers de logs
sous la forme suivante:
/index.html [19/05/2017:18:45:03 +0200]
/contact.html [19/05/2017:18:46:15 +0200]
/news.php?id=5 [24/05/2017:18:13:02 +0200]
/news.php?id=4 [24/05/2017:18:13:12 +0200]
/news.php?id=18 [24/05/2017:18:14:31 +0200]
...etc...
"
Ici, notre clef sera par exemple l'URL dacc s la page, et nos
op rations MAP et REDUCE seront exactement les m mes que celles
qui viennent d' tre pr sent es: on obtiendra ainsi le nombre de vue
pour chaque page distincte du site.
COURS BIG DATA - 2019
23
MapReduce : Exercice Graphe social (1/8)
" Un autre exemple: on administre un r seau social comportant des
millions d'utilisateurs.
" Pour chaque utilisateur, on a dans notre base de donn es la liste des
utilisateurs qui sont ses amis sur le r seau (via une requ te SQL).
" On souhaite afficher quand un utilisateur va sur la page d'un autre
utilisateur une indication Vous avez N amis en commun .
" On ne peut pas se permettre d'effectuer une s rie de requ tes SQL
chaque fois que la page est acc d e (trop lourd en traitement). On va
donc d velopper des programmes MAP et REDUCE pour cette
op ration et ex cuter le traitement toutes les nuits sur notre base de
donn es, en stockant le r sultat dans une nouvelle table.
COURS BIG DATA - 2019
24
MapReduce : Exercice Graphe social (2/8)
"
Ici, nos donn es d'entr e sous la forme Utilisateur =>Amis:
A => B, C, D
B => A, C, D, E
C => A, B, D, E
D => A, B, C, E
E => B, C, D
" Puisqu'on est int ress par l'information amis en commun entre deux
utilisateurs et qu'on aura terme une valeur par clef, on va choisir
pour clef la concat nation entre deux utilisateurs. Par exemple, la clef
A-B d signera les amis en communs des utilisateurs A et B .
" On peut segmenter les donn es d'entr e l aussi par ligne.
COURS BIG DATA - 2019
25
MapReduce : Exercice Graphe social (3/8)
" Notre op ration MAP va se contenter de prendre la liste des amis
fournie en entr e, et va g n rer toutes les clefs distinctes possibles
partir de cette liste. La valeur sera simplement la liste d'amis, telle
quelle.
" On fait galement en sorte que la clef soit toujours tri e par ordre
Advertisement
alphab tique (clef B-A sera exprim e sous la forme A-B ).
" Ce traitement peut para tre contre-intuitif, mais il va terme nous
permettre d'obtenir, pour
couples
chaque
(clef;valeur): les deux listes d'amis de chacun des utilisateurs qui
composent la clef.
clef distincte, deux
COURS BIG DATA - 2019
26
MapReduce : Exercice Graphe social (4/8)
" Le pseudo code de notre op ration MAP:
UTILISATEUR =
POUR AMI dans , FAIRE:
SI UTILISATEUR < AMI:
CLEF = UTILISATEUR+"-"+AMI
SINON:
CLEF = AMI+"-"+UTILISATEUR
GENERER COUPLE (CLEF; )
" Par exemple, pour la premi re ligne:
A => B, C, D
On obtiendra les couples (clef;valeur):
("A-B"; "B C D")
("A-C"; "B C D")
("A-D"; "B C D")
COURS BIG DATA - 2019
27
MapReduce : Exercice Graphe social (5/8)
" Pour la seconde ligne :
B => A, C, D, E
On obtiendra ainsi :
("A-B"; "A C D E")
("B-C"; "A C D E")
("B-D"; "A C D E")
("B-E"; " A C D E")
" Pour la troisi me ligne :
C => A, B, D, E
On aura :
("A-C"; "A B D E")
("B-C"; "A B D E")
("C-D"; "A B D E")
("C-E"; " A B D E")
"
...et ainsi de suite pour nos 5 lignes d'entr e.
COURS BIG DATA - 2019
28
MapReduce : Exercice Graphe social (6/8)
" Une fois l'op ration MAP effectu e, Hadoop va r cup rer les couples
(clef;valeur) de tous les fragments et les grouper par clef distincte. Le
r sultat sur la base de nos donn es d'entr e :
Pour la clef "A-B " : valeurs "A C D E" et "B C D"
Pour la clef "A-C " : valeurs "A B D E" et "B C D"
Pour la clef "A-D " : valeurs "A B C E" et "B C D"
Pour la clef "B-C " : valeurs "A B D E" et "A C D E"
Pour la clef "B-D " : valeurs "A B C E" et "A C D E"
Pour la clef "B-E " : valeurs "A C D E" et "B C D"
Pour la clef "C-D " : valeurs "A B C E" et "A B D E"
Pour la clef "C-E " : valeurs "A B D E" et "B C D"
Pour la clef "D-E " : valeurs "A B C E" et "B C D"
" & on obtient bien, pour chaque clef USER1-USER2 , deux listes
d'amis: les amis de USER1 et ceux de USER2.
COURS BIG DATA - 2019
29
MapReduce : Exercice Graphe social (7/8)
"
Il nous faut enfin crire notre programme REDUCE. Il va recevoir en
entr e toutes les valeurs associ es une clef. Son r le va tre tr s
simple: d terminer quels sont les amis qui apparaissent dans les listes
(les valeurs) qui nous sont fournies. Pseudo-code:
LISTE_AMIS_COMMUNS=[] // Liste vide au d part.
SI LONGUEUR(VALEURS)!=2, ALORS: // Ne devrait pas se produire.
RENVOYER ERREUR
SINON :
POUR AMI DANS VALEURS[0], FAIRE :
SI AMI EGALEMENT PRESENT DANS VALEURS[1], ALORS :
// Pr sent dans les deux listes d'amis, on l'ajoute.
LISTE_AMIS_COMMUNS+=AMI
RENVOYER LISTE_AMIS_COMMUNS
COURS BIG DATA - 2019
30
MapReduce : Exercice Graphe social (8/8)
" Apr s ex cution de l'op ration REDUCE pour les valeurs de chaque clef
unique, on obtiendra donc, pour une clef A-B , les utilisateurs qui
apparaissent dans la liste des amis de A et dans la liste des amis de B.
Autrement dit, on obtiendra la liste des amis en commun des utilisateurs
A et B. Le r sultat:
"A-B " : "C, D"
"A-C " : "B, D"
"A-D " : "B, C"
"B-C " : "A, D, E"
"B-D " : "A, C, E"
"B-E " : "C, D"
"C-D " : "A, B, E"
"C-E " : "B, D"
"D-E " : "B, C"
" On sait ainsi que A et B ont pour amis
communs les utilisateurs C et D, ou encore
que B et C ont pour amis communs les
utilisateurs A, D et E.
COURS BIG DATA - 2019
31
MapReduce : Conclusion
" En utilisant
le mod le MapReduce, on a ainsi pu cr er deux
programmes tr s simples (nos programmes MAP et REDUCE) de
quelques lignes de code seulement, qui permettent d'effectuer un
traitement somme toute assez complexe.
" Mieux encore, notre traitement est parall lisable: m me avec des
dizaines de millions d'utilisateurs, du moment qu'on a assez de
le traitement sera effectu
machines au sein du cluster Hadoop,
rapidement. Pour aller plus vite,
il nous suffit de rajouter plus de
machines.
" Pour notre r seau social, il suffira d'effectuer ce traitement toutes les
nuits heure fixe, et de stocker les r sultats dans une table. Ainsi,
lorsqu'un utilisateur visitera la page d'un autre utilisateur, un seul
SELECT dans la base de donn es suffira pour obtenir la liste des amis
en commun avec un poids en traitement tr s faible pour le serveur.
COURS BIG DATA - 2019
32
Architecture Hadoop: Pr sentation (1/3)
Comme pour HDFS, la gestion des t ches de Hadoop se base sur deux
serveurs (des daemons):
" Le JobTracker, qui va directement recevoir la t che ex cuter (un .jar
Java), ainsi que les donn es d'entr es (nom des fichiers stock s sur HDFS)
et le r pertoire o stocker les donn es de sortie (toujours sur HDFS). Il y
a un seul JobTracker sur une seule machine du cluster Hadoop. Le
JobTracker est en communication avec le NameNode de HDFS et sait
donc o sont les donn es.
" Le TaskTracker, qui est en communication constante avec le JobTracker et
va recevoir les op rations simples effectuer (MAP/REDUCE) ainsi que
les blocs de donn es correspondants (stock s sur HDFS). Il y a un
TaskTracker sur chaque machine du cluster.
COURS BIG DATA - 2019
33
Architecture Hadoop: Pr sentation (2/3)
" Comme le JobTracker est conscient de la position des donn es
(gr ce au NameNode), il peut facilement d terminer les meilleures
machines auxquelles attribuer les sous-t ches (celles o les blocs de
donn es correspondants sont stock s).
" Pour effectuer un traitement Hadoop, on va donc stocker nos
d'entr e sur HDFS, cr er un r pertoire o Hadoop
donn es
stockera les r sultats sur HDFS, et compiler nos programmes MAP
et REDUCE au sein d'un .jar Java.
" On soumettra alors le nom des fichiers d'entr e,
le nom du
r pertoire des r sultats, et le .jar lui-m me au JobTracker: il
s'occupera du reste (et notamment de transmettre les programmes
MAP et REDUCE aux
serveurs TaskTracker des machines du
cluster).
COURS BIG DATA - 2019
34
Architecture Hadoop: Pr sentation (3/3)
COURS BIG DATA - 2019
35
Architecture Hadoop: Le JobTracker (1/3)
Le d roulement de l'ex cution d'une t che Hadoop suit les tapes suivantes du
point de vue du JobTracker :
1.
2.
3.
4.
Le client (un outil Hadoop console) va soumettre le travail effectuer au JobTracker:
une archive java .jar impl mentant les op rations Map et Reduce. Il va galement
soumettre le nom des fichiers d'entr e et l'endroit o stocker les r sultats.
Le JobTracker communique avec le NameNode HDFS pour savoir o se trouvent les
blocs correspondant aux noms de fichiers donn s par le client.
les noeuds
Le JobTracker, partir de ces informations, d termine quels sont
TaskTracker les plus appropri s, c'est dire ceux qui contiennent les donn es sur
lesquelles travailler sur la m me machine, ou le plus proche possible (m me rack/rack
Advertisement
proche).
Pour chaque morceau des donn es d'entr e, le JobTracker envoie au TaskTracker
s lectionn le travail effectuer (MAP/REDUCE, code Java) et les blocs de donn es
correspondants.
COURS BIG DATA - 2019
36
Architecture Hadoop: Le JobTracker (2/3)
5.
6.
7.
Pour chaque morceau des donn es d'entr e, le JobTracker envoie au TaskTracker
s lectionn le travail effectuer (MAP/REDUCE, code Java) et les blocs de donn es
correspondants.
Le JobTracker communique avec les noeuds TaskTracker en train d'ex cuter les t ches. Ils
envoient r guli rement un heartbeat , un message signalant qu'ils travaillent toujours
sur la sous-t che re ue. Si aucun heartbeat n'est re u dans une
le
JobTracker consid re la t che comme ayant chou e et donne le m me travail effectuer
un autre TaskTracker.
p riode donn e,
Si par hasard une t che choue (erreur java, donn es incorrectes, etc.), le TaskTracker va
signaler au JobTracker que la t che n'a pas pu tre ex cut e.
Le JobTracker va alors d cider de la conduite adopter :
- Demander au m me TaskTracker de r -essayer.
- Redonner la sous-t che un autre TaskTracker.
- Marquer les donn es concern es comme invalides, etc.
Il pourra m me blacklister le TaskTracker concern comme non-fiable
49
Architecture Hadoop: Le JobTracker (3/3)
7. Une fois que toutes les op rations envoy es aux TaskTracker (MAP + REDUCE) ont
t effectu es et confirm es comme effectu es par tous les noeuds, le JobTracker
marque la t che comme effectu e . Des informations d taill es sont disponibles
(statistiques, TaskTracker ayant pos probl me, etc.).
Remarques
" Par ailleurs, on peut galement obtenir tout moment de la part du JobTracker
des informations sur les t ches en train d' tre effectu es: tape actuelle (MAP,
SHUFFLE, REDUCE), pourcentage de compl tion, etc.
" La soumission du .jar, l'obtention de ces informations, et d'une mani re g n rale
toutes les op rations li es Hadoop s'effectuent avec le m me unique client
console vu pr c demment: hadoop (avec d'autres options que l'option fs vu
pr c demment).
COURS BIG DATA - 2019
38
Architecture Hadoop: Le TaskTracker
" Le TaskTracker dispose d'un nombre de slots d'ex cution. A chaque slot
correspond une t che ex cutable (configurable). Ainsi, une machine ayant par
exemple un processeur 8 cSurs pourrait avoir 16 slots d'op rations configur es.
" Lorsqu'il re oit une nouvelle t che effectuer (MAP, REDUCE, SHUFFLE) depuis
le JobTracker, le TaskTracker va d marrer une nouvelle instance de Java avec le
fichier .jar fourni par le JobTracker, en appelant l'op ration correspondante.
" Une fois la t che d marr e, il enverra r guli rement au JobTracker ses messages
heartbeats. En dehors d'informer le JobTracker qu'il est toujours fonctionnels,
ces messages indiquent galement
le nombre de slots disponibles sur le
TaskTracker concern .
" Lorsqu'une sous-t che est termin e,
le TaskTracker envoie un message au
JobTracker pour l'en informer, que la t che se soit bien d roul e ou non (il
indique videmment le r sultat au JobTracker).
COURS BIG DATA - 2019
39
Architecture Hadoop: Remarques (1/2)
" De mani re similaire au NameNode de HDFS,
il n'y a qu'un seul
JobTracker et s'il tombe en panne, le cluster tout entier ne peut plus
effectuer de t ches. L aussi, des r solutions aux probl mes sont ajout es
dans la version 2 de Hadoop (explication dans la section suivante).
" G n ralement, on place le JobTracker et le NameNode HDFS sur la
m me machine (une machine plus puissante que les autres), sans y
placer de TaskTracker/DataNode HDFS pour limiter la charge. Cette
les deux
machine particuli re
gestionnaires , de t ches et de fichiers) est commun ment appel e le noeud
ma tre ( Master Node ). Les autres noeuds (contenant TaskTracker +
DataNode) sont commun ment appel s noeuds esclaves ( slave node ).
au sein du cluster
contient
(qui
COURS BIG DATA - 2019
40
Architecture Hadoop: Remarques (2/2)
" M me si le JobTracker est situ sur une seule machine, le client qui
envoie la t che au JobTracker initialement peut tre ex cut sur n'importe
quelle machine du cluster comme les TaskTracker sont pr sents sur la
machine, ils indiquent au client comment joindre le JobTracker.
" La m me remarque est valable pour l'acc s au syst me de fichiers: les
DataNodes indiquent au client comment acc der au NameNode.
" Enfin,
tout changement de configuration Hadoop peut s'effectuer
facilement simplement en changeant la configuration sur la machine o
sont situ s les serveurs NameNode et JobTracker: ils r pliquent les
changements de configuration sur tout le cluster automatiquement.
COURS BIG DATA - 2019
41
Architecture Hadoop: Architecture g n rale
COURS BIG DATA - 2019
42
YARN (MapReduce 2) : Pr sentation
" YARN (Yet-Another-Resource-Negotiator)
est aussi appel MRv2
(MapReduce 2). Ce nest 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 dallocation des ressources.
COURS BIG DATA - 2019
43
YARN (MapReduce 2) : Architecture (1/7)
Le JobTracker a trop de responsabilit s.
" G rer les ressources du cluster.
" G rer tous les jobs
JobTracker
Allouer les t ches et les ordonnancer.
Monitorer l'ex cution des t ches.
G rer le fail-over.
RessourceManager
ApplicationMaster
AM AM
Re-penser larchitecture du JobTracker.
"
S parer la gestion des ressources du cluster de la
coordination des jobs.
" Utiliser les noeuds esclaves pour g rer les jobs.
ResourceManager et ApplicationMaster.
" ResourceManager remplace le JobTracker et ne
g re que les ressources du Cluster.
" Une entit ApplicationMaster est allou e par
Application pour g rer les t ches.
" ApplicationMaster est d ploy e sur les noeuds
esclaves.
Cours Big Data - 2019
56
YARN (MapReduce 2) : Architecture (2/7)
" Cette nouvelle version contient aussi un autre composant :
Le NodeManager (NM)
Permet dex cuter plus de t ches qui ont du sens pour lApplication
Master, pas seulement du Map et du Reduce.
La taille des ressources est variable (RAM, CPU, network&.). Il y aura
plus de valeurs cod es en dur qui n cessitent un red marrage.
RessourceManager
AM AM ApplicationMaster
NodeManager
NM NM
COURS BIG DATA - 2019
45
YARN (MapReduce 2) : Architecture (3/7)
" Le JobTracker a disparu de larchitecture, ou plus pr cis ment, ses
r les ont t r partis diff remment.
" Larchitecture est maintenant organis e autour dun ResourceManager dont
le p rim tre daction est global au cluster et des ApplicationMaster
locaux dont le p rim tre est celui dun job ou dun groupe de jobs.
" En terme de responsabilit s, on peut donc dire que :
JobTracker = ResourceManager +ApplicationMaster.
" La diff rence, de part le d couplage, se trouve dans la multiplicit . En effet,
Un ResourceManager g re n ApplicationMaster, lesquels g rent chacun n
jobs.
COURS BIG DATA - 2019
46
YARN (MapReduce 2) : Architecture (4/7)
COURS BIG DATA - 2019
47
YARN (MapReduce 2) : Architecture (5/7)
1. ResourceManager
" Le ResourceManager est le rempla ant du JobTracker du point de vue du
client qui soumet des jobs (ou plut t des applications en Hadoop 2) un cluster
Hadoop.
Il na maintenant plus que deux t ches bien distinctes accomplir :
Scheduler
ApplicationsManager
"
a. Scheduler
" Le Scheduler est responsable de lallocation des ressources des applications
tournant sur le cluster.
Il sagit uniquement dordonnancement et dallocation de ressources.
Advertisement
"
" Les ressources allou es aux applications par le Scheduler pour leur permettre de
sex cuter sont appel es des Containers.
COURS BIG DATA - 2019
48
YARN (MapReduce 2) : Architecture (6/7)
v Container
" Un Container d signe un regroupement de m moire, de cpu, despace disque, de
bande passante r seau, &
b. ApplicationsManager
" LApplicationsManager accepte les soumissions dapplications.
" Une application n tant pas g r e par
la partie
ApplicationsManager ne soccupe que de n gocier le premier Container que le
Scheduler allouera sur un noeud du cluster. La particularit de ce premier
Container est quil contient lApplicationMaster (diapo suivant).
le ResourceManager,
2. NodeManager
" Les NodeManager sont des agents tournant sur chaque nSud et
le
Scheduler au fait de l volution des ressources disponibles. Ce dernier peut ainsi
prendre ses d cisions dallocation des Containers en prenant en compte des
demandes de ressources cpu, disque, r seau, m moire, &
tenant
Cours Big Data - 2019
61
YARN (MapReduce 2) : Architecture (7/7)
3. ApplicationMaster
" LApplicationMaster est le composant sp cifique chaque application, il est en
charge des jobs qui y sont associ s.
Lancer et au besoin relancer des jobs
N gocier les Containers n cessaires aupr s du Scheduler
Superviser l tat et la progression des jobs.
" Un ApplicationMaster g re donc un ou plusieurs jobs tournant sur un framework
donn . Dans le cas de base, cest donc un ApplicationMaster qui lance un job
MapReduce. De ce point de vue, il remplit un r le de TaskTracker.
v LApplicationsManager est lautorit qui g re les ApplicationMaster du
cluster. A ce titre, cest donc via lApplicationsManager que lon peut
Superviser l tat des ApplicationMaster
Relancer des ApplicationMaster
COURS BIG DATA - 2019
50
YARN (MapR2): D roulement de lex cution
"
Le d roulement de l'ex cution d'une t che Hadoop suit les tapes suivantes:
1.
2.
3.
4.
(un outil Hadoop console) va soumettre le travail effectuer au
Le client
ResourceManager: une archive java .jar impl mentant les op rations Map et Reduce, et
galement une classe driver (quon peut consid rer comme le main du programme).
Le ResourceManager alloue un container, Application Master , sur le cluster et y
lance la classe driver du programme.
Cet application master va se lancer et confirmer au ResourceManager quil tourne
correctement.
Pour chacun des fragments des donn es dentr e sur lesquelles travailler, lApplication
Master va demander au Resource Manager dallouer un container en lui indiquant les
donn es sur lesquelles celui-ci va travailler et le code ex cuter.
COURS BIG DATA - 2019
51
YARN (MapR2): D roulement de lex cution
5.
6.
7.
LApplication Master va alors lancer le code en question (une classe Java, g n rallement
map ou reduce) sur le container allou . Il communiquera avec la t che directement (via
un protocole potentiellement propre au programme lanc ),
le
ResourceManager.
sans passer par
Ces t ches vont r guli rement contacter lApplication Master du programme pour
transmettre des informations de progression, de statut, etc. parall lement, chacun des
NodeManager communique en permanence avec le ResourceManager pour lui indiquer
son statut en terme de ressources (containers lanc s, RAM, CPU), mais sans information
sp cifique aux t ches ex cut es.
Pendant lex cution du programme, le client peut tout moment contacter le programme
en cours dex cution; pour ce faire, il communique directement avec lApplication
Master du programme en contactant
le container correspondant sans passer par le
ResourceManager du cluster.
8. A lissue de lex cution du programme, lApplication Master sarr te; son container est
lib r et est nouveau disponible pour de futures t ches.
COURS BIG DATA - 2019
66
YARN (MapR2) : Lancement dune Application dans un
Cluster Yarn
COURS BIG DATA - 2019
53
YARN (MapR2) : Lancement dune Application dans un
Cluster Yarn
COURS BIG DATA - 2019
54
YARN (MapR2) : Lancement dune Application dans un
Cluster Yarn
COURS BIG DATA - 2019
55
YARN (MapR2) : Lancement dune Application dans un
Cluster Yarn
COURS BIG DATA - 2019
56
YARN (MapR2) : Lancement dune Application dans un
Cluster Yarn
COURS BIG DATA - 2019
57
YARN (MapR2) : Ex cution dun Job MR
COURS BIG DATA - 2019
58
YARN (MapR2) : Ex cution dun Job MR
COURS BIG DATA - 2019
59
YARN (MapR2) : Ex cution dun Job MR
COURS BIG DATA - 2019
60
YARN (MapR2) : Ex cution dun Job MR
COURS BIG DATA - 2019
61
YARN (MapR2) : Ex cution dun Job MR
COURS BIG DATA - 2019
62
Rappel: Hadoop 2: HDFS - Remarques
" De la m me fa on que l volution de MapReduce 2, nous trouvons quelques
am liorations de HDFS dans la 2 me version de Hadoop. Ces am liorations ont
t d j expliqu dans la partie consacr e HDFS.
1. Le SPOF du Namenode a disparu
"
Lancienne architecture dHadoop imposait lutilisation dun seul namenode, composant
central contenant les m tadonn es du cluster HDFS. Ce namenode tait un SPOF (Single
Point Of Failure).
" Dans Hadoop 2, on peut mettre deux namenodes en mode actif/attente. Si le Namenode
principal est indisponible, le Namenode secondaire prend sa place.
2. La f d ration HDFS
" Dans un cluster HDFS, un namenode correspond un espace de nommage (namespace).
Dans Lancienne architecture dHadoop, on ne pouvait utiliser quun namenode par
cluster. La f d ration HDFS permet de supporter plusieurs namenodes et donc plusieurs
namespace sur un m me cluster. (Exemple : un NameNode pourrait g rer tous les fichiers
sous /user, et un second NameNode peut g rer les fichiers sous /share.)
COURS BIG DATA - 2019
63
Hadoop 2: Architecture g n rale
COURS BIG DATA - 2019
64
Rappel: Hadoop 2: Une volution architecturale majeure
" Cette remise plat des r les permis aussi de d coupler Hadoop de
MapReduce et, ce faisant, de permettre des frameworks alternatifs d tre
port s directement sur Hadoop et HDFS.
" Cela va donc permettre Hadoop, outre une meilleure scalabilit , de
senrichir de nouveaux frameworks couvrant des besoins peu ou pas
couverts avec Map Reduce.
Cours Big Data - 2019
65
Rappel: Hadoop 2: Une volution architecturale majeure
Hadoop se transforme en OS de la donn e !
Client et cluster peuvent utiliser des versions diff rentes.
Des protocoles de communication standardis s et document s.
volution du framework progressive
avec
r tro-compatibilit
destruction des services.
Cours Big Data - 2019
sans
66
Hadoop 2 : cosyst me (1/2)
Cours Big Data - 2019
67
Hadoop 2 : cosyst me (2/2)
"
Pig : un langage de haut niveau d di l'analyse
de gros volumes de donn es. Il s'adresse aux
d veloppeurs habitu s faire des sc...