Cours de « Big Data »: HDFS et MapReduce

Page 1 sur 39Lecteur de document UniversityLib

Cours de « Big Data »: HDFS et MapReduce

Big Data and Distributed Systems · notes

Voir tous les documents en intelligence artificielle et données

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