MapReduce

Big Data · notes

Voir tous les documents en intelligence artificielle et données

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 : 3 Date Ville Produit Prix 2016-01-01 London Clothes 25.99 2016-01-01 Miami Music 12.15 2016-01-02 Miami Clothes 50.00 … … … … 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 Clé Valeur London 25,99 Miami 62,15 New York 3,10 4 Date Ville Produit Prix 2016-01-01 London Clothes 25.99 2016-01-01 Miami Music 12.15 2016-01-02 New York Toys 3,10 2016-01-02 Miami Clothes 50.00 … … … … 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 d’avoir 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 les données d'entrée en plusieurs fragments, et faire exécuter 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): (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 MapReduce : Exemple concret (3/8) Nos données d'entrée (le texte): … 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 MAPelle-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 l’occurrence 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 pseudocode): Pour chacun de nos fragments, les couples (clef; valeur) générés seront donc: 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 clefcommune. 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: 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: 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 l’occurrence, notre résultatsera: 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. L’intérêt du modèle MapReduce est qu'il nous suffit de développer les deux opérations réellement importantes du traitement: MAP et REDUCE, et de bénéficier automatiquement de la possibilité d'effectuer le traitement sur un nombre variable de machines de manière distribuée. 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: Ici, notre clef sera par exemple l'URL d’accè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êteSQL). 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: 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 parligne. 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 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 chaque clef distincte, deux couples (clef;valeur): les deux listes d'amis de chacun des utilisateurs qui composent la clef. COURS BIG DATA - 2019 26 MapReduce : Exercice – Graphe social (4/8) Le pseudo code de notre opération MAP:  On obtiendra les couples (clef;valeur): COURS BIG DATA - 2019 27 MapReduce : Exercice – Graphe social (5/8) Pour la seconde ligne :  On obtiendra ainsi : Pour la troisième ligne :  On aura : ...et ainsi de suite pour nos 5 lignesd'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 : … 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: 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: 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 machines au sein du cluster Hadoop, le traitement sera effectué 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 leserveur. 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 données d'entrée sur HDFS, créer un répertoire où Hadoop 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 : 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. Le JobTracker, à partir de ces informations, détermine quels sont les noeuds 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 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) 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 période donnée, le JobTracker considère la tâche comme ayant échouée et donne le même travail à effectuer à un autre TaskTracker. 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) 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 cœurs 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 machine particulière au sein du cluster (qui contient les deux «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 ») . 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 clusterautomatiquement. 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 n’est 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 d’allocation 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** - Allouer les tâches et les ordonnancer. - Monitorer l'exécution des tâches. - Gérer le fail-over. Re-penser l’architecture 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 etApplicationMaster. 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 d’exécuter plus de tâches qui ont du sens pour l’Application 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 unredémarrage. COURS BIG DATA - 2019 45 YARN (MapReduce 2) : Architecture (3/7) Le JobTracker a disparu de l’architecture, ou plus précisément, ses rôles ont été répartis différemment. L’architecture est maintenant organisée autour d’un ResourceManager dont le périmètre d’action est global au cluster et à des ApplicationMaster locaux dont le périmètre est celui d’un job ou d’un 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 n’a maintenant plus que deux tâches bien distinctes à accomplir : Scheduler ApplicationsManager a. Scheduler Le Scheduler est responsable de l’allocation des ressources des applications tournant sur le cluster. Il s’agit uniquement d’ordonnancement et d’allocation de ressources. Les ressources allouées aux applications par le Scheduler pour leur permettre de s’exécuter sont appelées des Containers . COURS BIG DATA - 2019 48 YARN (MapReduce 2) : Architecture (6/7)  Container Un Container désigne un regroupement de mémoire, de cpu, d’espace disque, de bande passante réseau, … b. ApplicationsManager L’ApplicationsManager accepte les soumissions d’applications. Une application n’étant pas gérée par le ResourceManager , la partie ApplicationsManager ne s’occupe que de négocier le premier Container que le Scheduler allouera sur un noeud du cluster. La particularité de ce premier Container est qu’il contient l ’ ApplicationMaster (diapo suivant) . 2. NodeManager Les NodeManager sont des agents tournant sur chaque nœud et tenant le Scheduler au fait de l’évolution des ressources disponibles. Ce dernier peut ainsi prendre ses décisions d’allocation des Containers en prenant en compte des demandes de ressources cpu, disque, réseau, mémoire, … Cours Big Data - 2019 61 YARN (MapReduce 2) : Architecture (7/7) 3. ApplicationMaster L’ApplicationMaster 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, c’est donc un ApplicationMaster qui lance un job MapReduce. De ce point de vue, il remplit un rôle de TaskTracker .  L’ApplicationsManager est l’autorité qui gère les ApplicationMaster du cluster . A ce titre, c’est donc via l’ApplicationsManager que l’on peut Superviser l’état desApplicationMaster Relancer desApplicationMaster COURS BIG DATA - 2019 50 YARN (MapR2): Déroulement de l’exécution Le déroulement de l'exécution d'une tâche Hadoop suit les étapes suivantes: Le client (un outil Hadoop console) va soumettre le travail à effectuer au ResourceManager: une archive java .jar implémentant les opérations Map et Reduce, et également une classe driver (qu’on 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 qu’il tourne correctement. Pour chacun des fragments des données d’entrée sur lesquelles travailler, l’Application Master va demander au Resource Manager d’allouer 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 l’exécution L’Application 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é), sans passer par le ResourceManager. Ces tâches vont régulièrement contacter l’Application 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 l’exécution du programme, le client peut à tout moment contacter le programme en cours d’exécution; pour ce faire, il communique directement avec l’Application Master du programme en contactant le container correspondant sans passer par le ResourceManager du cluster. A l’issue de l’exécution du programme, l’Application Master s’arrète; son container est libéré et est à nouveau disponible pour de futures tâches. COURS BIG DATA - 2019 66 YARN (MapR2) : Lancement d’uneApplication dans un ClusterYarn COURS BIG DATA - 2019 53 YARN (MapR2) : Lancement d’uneApplication dans un ClusterYarn COURS BIG DATA - 2019 54 YARN (MapR2) : Lancement d’uneApplication dans un ClusterYarn COURS BIG DATA - 2019 55 YARN (MapR2) : Lancement d’uneApplication dans un ClusterYarn COURS BIG DATA - 2019 56 YARN (MapR2) : Lancement d’uneApplication dans un ClusterYarn COURS BIG DATA - 2019 57 YARN (MapR2) : Exécution d’un Job MR COURS BIG DATA - 2019 58 YARN (MapR2) : Exécution d’un Job MR COURS BIG DATA - 2019 59 YARN (MapR2) : Exécution d’un Job MR COURS BIG DATA - 2019 60 YARN (MapR2) : Exécution d’un Job MR COURS BIG DATA - 2019 61 YARN (MapR2) : Exécution d’un 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 L’ancienne architecture d’Hadoop imposait l’utilisation d’un 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 L’ancienne architecture d’Hadoop, on ne pouvait utiliser qu’un 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 s’enrichir 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. sans Cours Big Data - 2019 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 scripts via Bash ou Python par exemple. Hive : un système d'entrepôt de données (Data Warehouse) pour Hadoop qui offre un langage de requête proche de SQL pour faciliter les agrégations, le requêtage ad-hoc et l’analyse de gros volumes de données stockés dans des systèmes de fichiers compatibles Hadoop. Hbase : un système de gestion de base de données NoSQL orienté colonnes. Il est distribué, scalable et dédié au big data avec un accès direct et une lecture/écriture temps réel. Sqoop : un outil conçu pour transférer efficacement une masse de données entre Apache Hadoop et un stockage de données structuré tel que les bases de données relationnelles . Mahout : un système d' apprentissage automatique et d'analyse de données. Il implémente des algorithmes de classification et de regroupement automatique (Machine Learning, DataMining). Flume : un service distribué, fiable et disponible pour collecter efficacement, agréger et déplacer une grande quantité de logs . Zookeeper : un service centralisé pour maintenir les configurations , la nomenclature , pour fournir une synchronisation distribuée et des services groupés. Oozie : un système de flux de travail (workflow) dont l'objectif est de simplifier la coordination et la séquence de différents traitements (programme Map/Reduce, scripts Pig,…). 80 Hadoop 2: Conclusion Les changements architecturaux d’Hadoop 2 sont, à minima, intéressants à deux titres : Du point de vue ops, haute disponibilité et la capacité de mutualiser l’infrastructure HDFS. Du point de vue dev et métier, le découplage par rapport à MapReduce va permettre de faire émerger d’autres framework. Hadoop n’est donc plus un moteur MapReduce, mais bien un moteur générique pour du calcul distribué. COURS BIG DATA - 2019 69 Exemple Word Count 76 Design Patterns - MapReduce Les designs patterns (Patrons de Conception) représentent les types de traitements Map-Reduce les plus utilisés avec Hadoop. Ils sont classés en trois catégories : Patrons de Filtrage (Filtering Patterns) Echantillonnage de données Listes des top-n Patrons de Récapitulation (Summarization Patterns) Trouver les min et les max Statistiques Indexes Patrons Structurels Combinaison de données relationnelles 77 Patrons de filtrage Ne modifient pas les données Trient, parmi les données présentes, lesquelles garder et lesquelles enlever. o Filtrage sur les mappers : utilise principalement les tâches Mappers o Filtrage Top-n o Filtrage dé-duplication 78 Patrons de filtrage Exemple de Filtrage sur les mappers : Cas d’étude : fichier contenant tous les posts des utilisateurs sur un forum Filtre : Retenir les posts les plus courts, contenant une seule phrase. Une phrase est un post qui ne contient aucune ponctuation de la forme: .!?, ou alors une seule à la fin. 79 Patrons de filtrage Exemple : Top-10 Trouver parmi les différents posts des forums, les 10 posts les plus longs Traitement Map-Reduce (de la même manière qu’une sélection sportive)  Chaque Mapper génère une liste Top-10  Le Reducer trouve les Top 10 globaux 80 Patrons de filtrage Exemple : Patron de dé-duplication On cherche des informations sur les profils des internautes ayant accédé à un site web, il est pratique de pouvoir éliminer de l’ensemble des traces toutes les réplications d’un même internaute.  Les Mappers écrivent des paires clé-valeur avec les enregistrements retenus utilisés en tant que clé  Le mécanisme de regroupement des paires clé-valeur du MapReduce va réaliser automatiquement la dé-duplication. 81 Patrons de récapitulation Permettent de donner une idée de haut niveau sur les données en question. On distingue deux types : Récapitulation (ou résumé) numérique :  Chercher des chiffres, des comptes (combien dispose-t-on d’un certain type d’entrées)  Min et Max  Premier et dernier  Moyenne … Index : tels que les index utilisés par google pour représenter les pages web 82 Récapitulations numériques Peuvent être : Le nombre de mots, enregistrements… Souvent : la clé = l’objet à compter, et la valeur = 1 (nombre d’occurrences) Min-Max / Premier-Dernier La moyenne La médiane Écart type … Exemple : Word Counter : calcul du nombre d’occurrences de chaque mot dans un texte. 83 Patrons de récapitulation Utilisation d’un mélangeur Possibilité d’utiliser un mélangeur (Combiner) entre les mappers et les reducers. Permettent de réaliser la réduction en local dans chacun des Mappers AVANT de faire appel au nœud réducteur principal. 84 Patrons de récapitulation Exemple : Calcul du prix de vente moyen de divers produits identifiés par leurs références ( Id ). 85 Patrons de récapitulation Sans Combiner o Les Mappers parcourent les données, et générent des couples (Id, prix) o Pour chaque id produit, le Reducer calcule la somme et incrémente un compteur o Il divise à la fin la somme par le compteur 86 Patrons de récapitulation Avec Combiner o Chaque nœud réalise une première réduction où les sommes locales sont calculées o Le Reducer final regroupe ces sommes et synthétise la moyenne finale o Nombre d’enregistrements envoyés au réducteur significativement réduit o Temps nécessaire pour la réduction diminue 87 Index Les index permettent une recherche plus rapide Dans un livre : pour chaque mot donné, indiquer les différentes pages où se trouve ce mot. Dans le web : on trouve des liens vers des pages web à partir d’un ensemble de mots-clés. Exemple : indexation des mots dans les posts d’un forum 88 Index Les Mappers se chargent de lire et d’analyser les fichiers d’entrée. Ils dressent la liste de tous les mots contenus dans chaque fichier Ils génèrent des paires clé-valeur constituées d’un mot reconnu (la clé), et de la référence du document (la valeur). Habituellement on réalise un Mapper intelligent, qui élimine certains mots considérés comme non-intéressants Le Shuffle & Sort est utile pour agréger les résultats, et associer chaque clé, donc chaque mot, à une liste de références de documents. Le Reducer n’a pas de traitement à faire dans sa fonction reduce . 89 Patrons structurels Utilisés quand les données proviennent d’une base de données structurée : Plusieurs tables, donc plusieurs sources de données, liées par clé étrangère Les données des différentes tables sont exportées sous forme de fichiers délimités 90 Patrons structurels Tâche du Mapper : Parcourir l’ensemble des fichiers correspondant aux tables de la base. Extraire de chacune des entrées les données nécessaires, en utilisant comme clé la clé étrangère joignant les deux tables Afficher les données extraites des différentes tables, chacune sur une ligne, en créant un champ supplémentaire indiquant la source des données. 91 Patrons structurels Tâche du Reducer : Faire l’opération de jointure entre les deux sources, en testant la provenance avec le champ supplémentaire (A ou B) 92 Conclusion - Avantages Gestion des défaillances : - Que ce soit au niveau du stockage ou du traitement, les nœuds responsables de ces opérations sont automatiquement gérés en cas de défaillance. Sécurité et persistance des données : - Grâce au concept « Rack Awarness », il n’y a plus de soucis de perte de données. Montée en charge : o Garantie d’une montée en charge maximale. Complexité réduite : Capacité d'analyse et de traitement des données à grande échelle. Coût réduit : o Hadoop est open source, et malgré leur massivité et complexité, les données sont traitées efficacement et à très faible coût 93 Conclusion - Inconvénients Difficulté d’intégration avec d’autres systèmes informatiques : Le transfert de données d’une structure Hadoop vers des bases de données traditionnelles est loin d’être trivial Administration complexe : Hadoop utilise son propre langage. L’entreprise doit donc développer une expertise spécifique Hadoop ou faire appel à des prestataires extérieurs Traitement de données différé et temps de latence important : Hadoop n’est pas fait pour l’analyse temps réel des données. Produit en développement continu Manque de maturité 94