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
DataNode
• interieurement, 1 block HDFS est soutenu par
multiple operating system blocks
7
HDFS: Blocks
Publicité
8
HDFS: Replcation
• Les blocs de données sont répliquées sur
plusieurs nœuds – 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
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)
Publicité
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 l’exemple 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
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
Publicité
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 à 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 s’occupe d’une 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 d’autres tâches.
38
Fault tolerance
39