Cours de Big Data "
HDFS MapReduce
Tol rance aux pannes
1
HDFS
" HDFS stocke les donn es sur plusieurs nodes
" Si HDFS d tecte un node d fectueux, alors il r alise la
fiabilit par la r plication des donn es sur plusieurs
autres nodes.
" Le syst me de fichiers est construit partir d'un groupe
de nodes de donn es , dont chacun g re des blocs de
donn es sur le r seau en utilisant un protocole de bloc
sp cifique HDFS .
2
HDFS
3
HDFS: Overview (from IBM)
4
Architecture HDFS
" L Architecture Ma tre / Esclave
" Master: NameNode
G re le syst me de fichiers espace de noms et les
m tadonn es
" FsImage
" EditLog
R glemente l'acc s des clients aux fichiers
" Slave: DataNode
Beaucoup par cluster
G re le stockage attach aux nodes
Rapporte p riodiquement l tat (Status) au NameNode
5
Architecture HDFS (By IBM)
6
HDFS: Blocks
" HDFS est con u pour supporter de tr s larges
fichiers
" Chaque fichier est divis en blocs: - Hadoop
d faut: 64Mo (mais - BigInsights d faut:
128Mo)
" Blocks r side physiquement sur diff rents
Publicité
DataNode
" interieurement, 1 block HDFS est soutenu par
multiple operating system blocks
7
HDFS: Blocks
8
HDFS: Replcation
" Les blocs de donn es sont r pliqu es sur
plusieurs nSuds
Le comportement est contr l par le facteur de
r plication, configurable par fichier
Par d faut est 3 r pliques
9
NameNode Startup
" 1. NameNode lit fsimage en m moire
" 2. NameNode applique les changements de editlog
" 3. NameNode attend les blocks data de Data nodes
NameNode ne stocke pas les informations de bloc
NameNode exprime un safemode lorsque 99,9%
des blocs ont au moins une copie comptabilis e
10
NameNode Startup
11
Ajout de fichier
" 1. Le fichier est ajout la m moire
NameNode et a persist dans editlog
" 2. Les donn es sont crites dans des blocs
DataNodes
DataNode commence une copie encha n e par
deux autres DataNodes
Si au moins une criture pour chaque bloc r ussit,
criture est r ussie
12
Ajout de fichier
13
FS - Shell sous Hadoop
14
FS - Shell sous Hadoop
15
Publicité
FS - Shell sous Hadoop
16
MapReduce Engine
Architecture Ma tre / Esclave
Un Ma tre unique (JobTracker) contr le l'ex cution du travail
sur plusieurs esclaves (TaskTrackers)
JobTracker
Accepte les travaux MapReduce soumis par les clients
Pushes map and reduce tasks out to TaskTracker nodes
Maintient le travail physiquement proche des donn es
Contr le les t ches et l tat de chaque TaskTracker
TaskTracker
Ex cute les t ches Map et reduce
Rapporte le statut JobTracker
G re le stockage et la transmission de sortie interm diaire
17
MapReduce Engine
18
MapReduce Programming Model
19
MapReduce Overview
20
MapReduce Overview
21
MAP
22
SHUFFLE
23
REDUCE
24
How does Hadoop run MapReduce
jobs?
25
Map Reduce et application
" Map Reduce
" Dans lexemple du compteur de mots, nous
avons une liste de noms d'animaux
" MapReduce peut automatiquement diviser les
fichiers sur les sauts de ligne
" Notre fichier a t divis en deux blocks sur
Publicité
deux Nodes
" Nous voulons compter combien de fois les
gros chats sont mentionn s.
Dans SQL qui serait:
26
Map Reduce et application
27
Map Input
" La fonction Map a besoin
du couple <Cl , Valeur> en
entr e
" Si aucune cl est
disponible, elle doit tre
fabriqu e.
" La mise en
correspondance entre
l'entr e (fichiers, lien
internet, ...) <Key, value>
se fait dans la classe
InputFormat
28
Map Input
29
Map Tasks
Nous avons deux t ches
" Filtrer les lignes qui ne
correspondent pas aux
gros chats
" Pr parer un comptage en
transformant <Cl ,
Valeur> <texte (nom),
Entier (1)>
30
Map Tasks
31
Shuffle
Le Shuffle d place toutes les
valeurs d'une cl du m me
Node
" La distribution se fait
Publicité
travers une classe Partitioner
" Reduce task peut tre
ex cut e sur des Nodes
arbitraires, dans notre
exemple Node 1 et 3
32
Shuffle
NB: Les r sultats sont stock s dans des blocks HDFS
sur les machines qui ex cutent le Reduce Job
33
Reduce
" La fonction Reduce calcule les valeurs agr g es
pour chaque cl . Normalement, la sortie est crite
dans un DFS.
34
Reduce
35
Fault tolerance
36
Fault tolerance
" Si une t che enfant choue, l'enfant JVM raporte la d faillance
au TaskTracker et sort. La Tentative est marqu e chou , il
soccupe dune autre t che.
" Si la t che de l'enfant se bloque, il est tu . JobTracker replanifie
la t che sur une autre machine.
" Si la t che continue chouer, le travail a chou .
37
Fault tolerance
" TaskTracker chec
" JobTracker ne re oit aucun battement de
coeur
" Supprime TaskTracker pour planifier
dautres t ches.
38
Fault tolerance
39