Chapitre 2
Hadoop
Cours Big Data - 2019
1
Hadoop
Hadoop (High-Availability Distributed Oject-
Oriented Platform) :
Framework libre et open source
Ecrit en Java
G r par Apache
Capable de sattaquer un gros volume de
donn es de types diff rents
2
Hadoop: Pr sentation (1/3)
" Hadoop est une plateforme (framework) open source con ue pour
r aliser d'une fa on distribu e des traitements sur des volumes de
donn es massives, de l'ordre de plusieurs p taoctets. Ainsi, il est
destin
et
chelonnables (scalables). Il s'inscrit donc typiquement sur le terrain
du Big Data.
cr ation d'applications distribu es
faciliter
la
" Hadoop est g r sous l gide de la fondation Apache, il est crit en
Java.
" Hadoop a t con u par Doug Cutting en 2004 et t inspir par les
publications MapReduce, GoogleFS et BigTable de Google.
.
Cours Big Data - 2019
3
Hadoop: Pr sentation (2/3)
" Mod le simple pour les d veloppeurs: il suffit de d velopper des
t ches Map-Reduce, depuis des interfaces simples accessibles via
des librairies dans des langages multiples (Java, Python, C/C++...).
" D ployable
tr s
facilement
(paquets Linux
pr -configur s),
configuration tr s simple elle aussi.
" S'occupe de toutes les probl matiques li es au calcul distribu ,
comme lacc s et le partage des donn es, la tol rance aux pannes,
ou encore la r partition des t ches aux machines membres du cluster
: le programmeur a simplement s'occuper du d veloppement
logiciel pour l'ex cution de la t che.
Cours Big Data - 2019
4
Hadoop: Pr sentation (3/3)
Hadoop assure les crit res des solutions Big Data:
" Performance: support du traitement d' normes data sets (millions
de fichiers, Go To de donn es totales) en exploitant le parall lisme
et la m moire au sein des clusters de calcul.
" Economie: contr le des co ts en utilisant de mat riels de calcul de
type standard.
" Evolutivit (scalabilit ): un plus grand cluster devrait donner une
meilleure performance.
" Tol rance aux pannes: la d faillance d'un nSud ne provoque pas
l' chec de calcul.
" Parall lisme de donn es: le m me calcul effectu sur toutes les
donn es.
Cours Big Data - 2019
5
Histoire de Hadoop
2002
2004
2006
Doug Cutting qui travaille chez Google, tait lorigine
du projet Apache Lucence consistant d velopper une
librairie logicielle dindexation des recherches.
Apache Nutch a vu le jour comme un moteur de
recherche pour lensemble dInternet.
Au m me moment Google a publi : Google File
System ainsi que le traitement MapReduce.
Doug Cutting reprend les deux derniers concepts pour
aller plus loin dans les deux projets.
6
Histoire de Hadoop
2006
2008
2011
Hadoop sous-projet de Apache Lucene
Hadoop un projet ind pendant de la fondation Apache
Doug Cutting a rejoint Yahoo et a lanc , le premier
grand projet Hadoop, le Yahoo! Search Webmap qui
tourne sur un cluster de 10.000 cSurs Linux.
Yahoo rend le code source dHadoop public.
Doug Cutting a rejoint Cloudera, une soci t cr e par
deux de ses amis, galement contributeurs la
communaut Hadoop.
7
Hadoop
Plusieurs entreprises utilisent Hadoop :
Publicité
Amazon
Adobe
Ebay
Yahoo
Liste :
https://wiki.apache.org/hadoop/PoweredBy
10
Hadoop
r pond aux nouvelles
Hadoop fournit :
Syst me distribu qui
dimensions des donn es.
Syst me de fichiers HDFS (Hadoop Distributed
File System) permettant de stocker la donn e en la
dupliquant.
Syst me de traitement des gros volumes de
donn es appel MapReduce.
11
Hadoop
Principe :
Diviser les donn es
Les sauvegarder sur une collection de machines,
appel es cluster
Traiter les donn es directement l o elles sont
stock es, plut t que de les copier partir dun
serveur distribu
Il est possible dajouter des machines un
cluster, au fur et mesure que les donn es
augmentent.
Distributions de Hadoop
Trois principales distributions Hadoop :
Cloudera
Hortonworks
MapR
AWS (RosettaHub)
Echosyst me de Hadoop
En plus des briques de base Yarn/MapReduce/HDFS,
plusieurs outils existent pour permettre :
Lextraction et le stockage des donn es de/sur HDFS
La simplification des traitements sur ces donn es
La gestion et la coordination de la plateforme
Le monitoring du cluster
Hadoop : Choix
Cours Big Data - 2019
15
Echosyst me de Hadoop
Echosyst me de Hadoop
Echosyst me de Hadoop
Certains de ces outils se trouvent au dessus de la
couche Yarn/MR, tels que :
Pig : Langage de script
Hive : Langage proche de SQL (Hive
QL)
R Connectors : permet lacc s HDFS et
lex cution de requ tes Map/Reduce
partir du langage R
Mahout
learning et math matiques
Oozie : permet dordonnancer les jobs
Map Reduce, en d finissant des workflows
: biblioth que de machine
Echosyst me de Hadoop
Dautres outils se situent directement au dessus de
HDFS, tels que :
Hbase : Base de donn es NoSQL
orient e colonnes
Impala (ne se trouve pas dans la
figure) : permet le requ tage de donn es
directement partir de HDFS (ou de
Hbase) en utilisant des requ tes Hive
SQL
19
Echosyst me de Hadoop
Certains outils permettent de connecter HDFS aux
sources externes, tels que :
Sqoop : Lecture et criture des donn es
partir de bases de donn es externes
Flume : Collecte de logs et stockage
dans HDFS
20
Echosyst me de Hadoop
Enfin, dautres outils permettent
administration de Hadoop, tels que :
la gestion et
:
outil
la
Ambari
provisionnement,
monitoring des clusters.
Zookeeper
Publicité
centralis
informations
nommage
distribu e.
:
pour
de
de
et
fournit
pour
gestion
et
le
le
un
maintenir
configuration,
service
les
de
synchronisation
21
HDFS Hadoop Distributed File System
PC Lame
Datacenter Google
HDFS Hadoop Distributed File System
HDFS est un syst me de fichiers distribu , extensible
et portable
Ecrit en Java
Permet de stocker de tr s gros volumes de donn es sur
un Cluster.
Un cluster est un grand nombre de machines (nSuds)
quip es de disques durs et connect es entre elles.
23
Hadoop: Composants fondamentaux
" Hadoop est constitu de deux grandes parties :
Hadoop Distibuted File System HDFS : destin pour le stockage distribu
des donn es
Distributed Programing Framework - MapReduce : destin pour le traitement
distribu des donn es.
HADOOP
HDFS
MapReduce
Cours Big Data - 2019
24
HDFS : Pr sentation (1/3)
" HDFS (Hadoop Distributed File System) est un syst me de fichiers :
Distribu
Extensible
Portable
D velopp par Hadoop, inspir de GoogleFS et crit en Java.
Tol rant aux pannes
Con u pour stocker de tr s gros volumes de donn es sur un grand nombre
de machines peu couteuses quip es de disques durs banalis s.
" HDFS permet l'abstraction de l'architecture physique de stockage, afin de
manipuler un syst me de fichiers distribu comme s'il s'agissait d'un disque dur
unique.
" HDFS reprend de nombreux concepts propos s par des syst mes de fichiers
classiques comme ext2 pour Linux ou FAT pour Windows. Nous retrouvons
donc la notion de blocs (la plus petite unit que l'unit de stockage peut g rer),
les m tadonn es qui permettent de retrouver les blocs partir d'un nom de
fichier, les droits ou encore l'arborescence des r pertoires.
Cours Big Data - 2019
25
HDFS : Pr sentation (2/3)
" Toutefois, HDFS se d marque d'un syst me de fichiers classique pour les
principales raisons suivantes :
HDFS n'est pas solidaire du noyau du syst me d'exploitation. Il assure une
portabilit et peut tre d ploy sur diff rents syst mes d'exploitation. Un des
inconv nients est de devoir solliciter une application externe pour monter
une unit de disque HDFS.
HDFS est un syst me distribu . Sur un syst me classique, la taille du
disque est g n ralement consid r e comme la limite globale d'utilisation.
Dans un syst me distribu comme HDFS, chaque nSud d'un cluster
correspond un sous-ensemble du volume global de donn es du cluster.
Pour augmenter ce volume global,
il suffira d'ajouter de nouveaux
nSuds. On retrouvera galement dans HDFS, un service central appel
Namenode qui aura la t che de g rer les m tadonn es.
Cours Big Data - 2019
26
HDFS : Pr sentation (3/3)
HDFS utilise des tailles de blocs largement sup rieures ceux des
syst mes classiques. Par d faut, la taille est fix e 64 Mo. Il est toutefois
possible de monter 128 Mo, 256 Mo, 512 Mo voire 1 Go. Alors que sur
des syst mes classiques, la taille est g n ralement de 4 Ko, l'int r t de
fournir des tailles plus grandes permet de r duire le temps d'acc s un
bloc. Notez que si la taille du fichier est inf rieure la taille d'un bloc, le
fichier n'occupera pas la taille totale de ce bloc.
HDFS fournit un syst me de r plication des blocs dont le nombre de
la phase d' criture, chaque bloc
r plications est configurable. Pendant
correspondant au fichier est r pliqu sur plusieurs nSuds. Pour la phase de
Publicité
lecture, si un bloc est indisponible sur un nSud, des copies de ce bloc seront
disponibles sur d'autres nSuds.
Cours Big Data - 2019
27
HDFS Hadoop Distributed File System
Quand un fichier mydata.txt est enregistr dans HDFS,
il est d compos en grands blocs (par d faut 64Mo), chaque
bloc ayant un nom unique : blk_1, blk_2&
Les blocs sont num rot s et chaque fichier sait quels blocs il
occupe.
Les blocs dun m me fichier ne sont pas forc ment tous sur
la m me machine.
La r partition dun fichier sur plusieurs machines permet dy
acc der simultan ment par plusieurs processus.
HDFS Hadoop Distributed File System
Un cluster HDFS est constitu de machines jouant
diff rents r les exclusifs entre eux :
Ma tre HDFS ou namenode : contient tous les noms et blocs des
fichiers comme un gros annuaire t l phonique.
Datanodes :
contenu des fichiers.
toutes les autres machines stockent
les blocs du
DN
DN
DN
DN
NN
29
HDFS Hadoop Distributed File System
HDFS Hadoop Distributed File System
Probl mes possibles
Panne dun datanode
Panne du namenode
Panne r seau
DN
DN
DN
DN
NN
HDFS : Architecture (1/3)
" Une architecture de machines HDFS (aussi appel e cluster HDFS) repose sur
trois types de composants majeurs :
NameNode (noeud de nom) :
Un Namenode est un service central (g n ralement
appel aussi ma tre) qui s'occupe de g rer l' tat du
syst me de fichiers. Il maintient l'arborescence du
syst me de fichiers et les m tadonn es de l'ensemble
des fichiers et r pertoires d'un syst me Hadoop. Le
Namenode a une connaissance des Datanodes ( tudi s
juste apr s) dans lesquels les blocs sont stock s.
Cours Big Data - 2019
32
HDFS : Architecture (2/3)
Secondary Namenode :
Le Namenode dans l'architecture Hadoop est un point unique de d faillance (Single
Point of Failure en anglais). Si ce service est arr t , il n'y a pas un moyen de pouvoir
extraire les blocs d'un fichier donn . Pour r pondre cette probl matique, un Namenode
secondaire appel Secondary Namenode a t mis en place dans l'architecture Hadoop
version2. Son fonctionnement est relativement simple puisque le Namenode secondaire
v rifie p riodiquement l' tat du Namenode principal et copie les m tadonn es. Si le
Namenode principal est indisponible, le Namenode secondaire prend sa place.
Datanode :
Pr c demment, nous avons vu qu'un Datanode contient les blocs de donn es. En effet, il
stocke les blocs de donn es lui-m mes. Il y a un DataNode pour chaque machine au sein
du cluster. Les Datanodes sont sous les ordres du Namenode et sont surnomm s les
Workers. Ils sont donc sollicit s par les Namenodes lors des op rations de lecture et
d' criture.
Cours Big Data - 2019
33
HDFS : Architecture (3/3)
" Remarque : En plus de lajout du Namenode secondaire, nous trouvons dautres
am liorations de HDFS dans la 2 me version de Hadoop :
Les NameNodes tant le point unique pour le stockage et la gestion des m tadonn es, ils
peuvent tre un goulot d' tranglement pour soutenir un grand nombre de fichiers,
notamment lorsque ceux-ci sont de petite taille. En acceptant des espaces de noms
multiples desservis par des NameNodes s par s, le HDFS limite ce probl me. Par
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
34
HDFS Hadoop Distributed File System
Cas de panne de r seau
Donn es temporairement inaccessibles si cest une panne
r seau
HDFS Hadoop Distributed File System
Cas de panne dun datanode
Perte de donn es
Solutions :
Hadoop r plique chaque bloc 3 fois (par d faut)
Il choisit 3 nSuds au hasard, et place une copie du bloc dans
chacun deux
Si le nSud est en panne, le namenode le d tecte, et soccupe de
r pliquer encore les blocs qui y taient h berg s pour avoir
Publicité
toujours 3 copies.
Concept de Rack Awareness
HDFS Hadoop Distributed File System
Cas de panne dun namenode
Donn es perdues jamais si le namenode est d faillant = Mort du HDFS !
Solutions :
D finition dun namenode de secours appel secondary namenode ;
Il enregistre des sauvegardes du namenode intervalles r guliers ;
Il permet de reprendre le travail si le Namenode actif est d faillant.
37
HDFS Hadoop Distributed File System
Le namenode est vital pour HDFS mais unique.
Hadoop 2.x a introduit une nouvelle configuration appel e
high availability :
2 namenodes de secours se comportent comme des clones ;
Ils sont en tat dattente et mis jour en permanence
laide des services appel s JournalNodes ;
Les namenodes de secours sont capables de prendre le relai
instantan ment en cas de panne du namenode ;
Le secondary namenode fait
namenodes de secours ; il devient alors inutile.
la m me chose que les
HDFS: criture d'un fichier (1/2)
" Si on souhaite crire un fichier au sein de HDFS, on va utiliser la commande
principale de gestion de Hadoop: hadoop, avec l'option fs. Mettons qu'on
souhaite stocker le fichier page_livre.txt sur HDFS.
" Le programme va diviser le fichier en blocs de 64 Mo (ou autre, selon la
configuration) supposons qu'on ait ici 2 blocs. Il va ensuite annoncer au
NameNode: Je souhaite stocker ce fichier au sein de HDFS, sous le nom
page_livre.txt .
" Le NameNode va alors indiquer au programme qu'il doit stocker le bloc 1
sur le DataNode num ro 3, et le bloc 2 sur le DataNode num ro 1.
" Le client hadoop va alors contacter directement les DataNodes concern s et
les
leur demander de stocker les deux blocs en question. Par ailleurs,
DataNodes s'occuperont en informant le NameNode de r pliquer les
donn es entre eux pour viter toute perte de donn es.
Cours Big Data - 2019
42
HDFS: criture d'un fichier (2/2)
Cours Big Data - 2019
43
HDFS: Lecture d'un fichier (1/2)
" Si on souhaite lire un fichier au sein de HDFS, on utilise l aussi le client
Hadoop. Mettons qu'on souhaite lire le fichier page_livre.txt.
" Le client va contacter le NameNode, et lui indiquer Je souhaite lire le
fichier page_livre.txt . Le NameNode lui r pondra par exemple Il est
compos de deux blocs. Le premier est disponible sur le DataNode 3 et 2, le
second sur le DataNode 1 et 3 .
" L aussi,
le programme contactera les DataNodes directement et
leur
demandera de lui
transmettre les blocs concern s. En cas d'erreur/non
r ponse d'un des DataNode, il passe au suivant dans la liste fournie par le
NameNode.
Cours Big Data - 2019
44
HDFS: Lecture d'un fichier (2/2)
Cours Big Data - 2019
45
HDFS: La commande Hadoop fs
" Comme indiqu plus haut, la commande permettant de stocker ou extraire des
l'utilitaire console hadoop, avec l'option fs. Il r plique
fichiers de HDFS est
globalement les commandes syst mes standard Linux, et est tr s simple utiliser:
hadoop fs -put livre.txt /data_input/livre.txt
"
Pour stocker le fichier livre.txt sur HDFS dans le r pertoire /data_input.
"
Pour obtenir le fichier /data_input/livre.txt de HDFS et le stocker dans le fichier local
hadoop fs -get /data_input/livre.txt livre.txt
livre.txt.
"
hadoop fs -mkdir /data_input
Pour cr er le r pertoire /data_input
"
hadoop fs -rm /data_input/livre.txt
Pour supprimer le fichier /data_input/livre.txt
" D'autres commandes usuelles: -ls, -cp, -rmr, -du, etc...
Cours Big Data - 2019
46
Sources
" Cours
Ce chapitre est pris essentiellement du cours : Hadoop / Big Data - MBDS
- Benjamin Renaut (avec quelques modifications).
" Articles
Hadoop et le "Big Data , La Technologie Hadoop au coeur des projets
"Big Data" - St phane Goumard Tech day Hadoop, Spark - Arrow-
Institute.
Comparison Between Hadoop 2.x vs Hadoop 3.x Data Flair.
Cours Big Data - 2019
88