MapReduce

Programming, Big Data · course

Browse all intelligence artificielle et données documents

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...