MapReduce - Concepts, Examples, and Applications

Cette conférence présente le paradigme MapReduce, un modèle de programmation distribué permettant de traiter efficacement de très grands volumes de données. Elle s'inscrit dans un cours sur le Big Data et détaille les concepts fondamentaux, des exemples concrets d'application, ainsi que l'architecture Hadoop et son évolution vers YARN.

D'après le document MapReduce - Concepts, Examples, and Applications

Cet article a été rédigé automatiquement à partir du document source, puis vérifié avant publication.

MapReduce - Concepts, Examples, and Applications

Document source

MapReduce - Concepts, Examples, and Applications

Distributed Computing and Big Data Processing · PDF · 88 pages · 2016

Afficher l'aperçu du document

Consulter le document original →

Cette conférence présente le paradigme mapreduce-6e20911580">MapReduce, un modèle de programmation distribué permettant de traiter efficacement de très grands volumes de données. Elle s'inscrit dans un cours sur le Big Data et détaille les concepts fondamentaux, des exemples concrets d'application, ainsi que l'architecture Hadoop et son évolution vers YARN.

Introduction au paradigme MapReduce

MapReduce est une méthode conçue pour résoudre des problèmes de grande taille en les divisant en sous-problèmes plus petits, exécutables en parallèle sur un cluster de machines. Cette approche, dite « diviser pour régner », permet d'exploiter la puissance de calcul distribuée. Le paradigme MapReduce formalise deux opérations principales :

  • MAP : transforme les données d'entrée en une série de couples (clef, valeur). Cette opération est parallélisable, chaque machine traitant un fragment distinct des données.
  • REDUCE : traite toutes les valeurs associées à chaque clef distincte produite par l'opération MAP, regroupant les résultats pour chaque clef.

Le traitement MapReduce se déroule en quatre étapes : découpage (split) des données, mapping, regroupement (shuffle) par clef, puis réduction (reduce) des groupes. Ce modèle permet une exécution distribuée et efficace sur des volumes massifs de données.

Exemple concret : comptage des mots dans un texte

Pour illustrer MapReduce, considérons un texte en français dont on souhaite connaître la fréquence d'apparition des mots. Les étapes sont :

  • Découpage : le texte est divisé en fragments, par exemple ligne par ligne.
  • Prétraitement : suppression de la ponctuation, des accents, et conversion en minuscules pour homogénéiser les données.
  • Opération MAP : chaque fragment est parcouru mot par mot, générant pour chaque mot le couple (mot ; 1), indiquant une occurrence.
  • Shuffle : Hadoop regroupe automatiquement tous les couples par clef (mot), formant des groupes de valeurs associées.
  • Opération REDUCE : pour chaque mot, on additionne toutes les occurrences (valeurs) pour obtenir le total d'apparitions.

Par exemple, après traitement, le mot « qui » apparaît 4 fois, « celui », « croyait » et « fou » 2 fois chacun. Ce modèle trivial peut être étendu à des corpus beaucoup plus vastes, comme l'ensemble des textes d'une bibliothèque.

Autres exemples d'application de MapReduce

Statistiques web

Un autre cas d'usage est le comptage du nombre de visiteurs par page d'un site web à partir des fichiers de logs. La clef choisie est l'URL de la page, et les opérations MAP et REDUCE sont similaires à l'exemple précédent, permettant d'obtenir le nombre de vues par page.

Calcul d'amis en commun dans un réseau social

Dans un réseau social avec des millions d'utilisateurs, on souhaite afficher, lors de la visite d'une page utilisateur, le nombre d'amis en commun avec le visiteur. Effectuer cette requête en temps réel serait trop coûteux. On utilise donc MapReduce pour pré-calculer ces données :

  • Données d'entrée : pour chaque utilisateur, la liste de ses amis.
  • Clef : la concaténation alphabétique de deux utilisateurs, par exemple « A-B », représentant la paire d'utilisateurs.
  • Opération MAP : pour chaque utilisateur et sa liste d'amis, on génère tous les couples possibles (clef ; liste d'amis). La clef est toujours ordonnée alphabétiquement.
  • Shuffle : Hadoop regroupe les listes d'amis associées à chaque clef (paire d'utilisateurs).
  • Opération REDUCE : pour chaque clef, on calcule l'intersection des deux listes d'amis, obtenant ainsi les amis communs.

Par exemple, pour la clef « A-B », on obtient la liste des amis communs « C, D ». Ce traitement est parallélisable et peut être exécuté régulièrement pour mettre à jour les données, permettant une consultation rapide lors de l'accès aux pages utilisateurs.

Architecture Hadoop pour MapReduce

Hadoop, un framework open source, implémente MapReduce et gère le stockage distribué via HDFS. L'architecture Hadoop repose sur deux types de serveurs :

  • JobTracker : serveur unique qui reçoit les tâches (fichiers .jar Java), connaît la localisation des données via le NameNode HDFS, et répartit les sous-tâches aux TaskTrackers.
  • TaskTracker : présent sur chaque machine du cluster, il exécute les opérations MAP et REDUCE sur les fragments de données locaux, et communique régulièrement son état au JobTracker.

Le JobTracker attribue les tâches aux TaskTrackers en privilégiant la localisation des données pour minimiser les transferts réseau. Il supervise également la progression des tâches et relance celles qui échouent, pouvant blacklister des nœuds défaillants.

Chaque TaskTracker dispose d'un nombre configurable de « slots » d'exécution, correspondant à des tâches parallèles pouvant être lancées simultanément. Lorsqu'une tâche est terminée, le TaskTracker informe le JobTracker du résultat.

Le JobTracker et le NameNode sont généralement hébergés sur une machine dédiée, appelée nœud maître, tandis que les autres machines sont des nœuds esclaves hébergeant les TaskTrackers et DataNodes.

Évolution vers YARN (MapReduce 2)

YARN (Yet Another Resource Negotiator), aussi appelé MapReduce 2, est une évolution majeure de l'architecture Hadoop visant à résoudre les limites de scalabilité et de gestion des ressources du JobTracker unique. YARN sépare la gestion des ressources du cluster de la coordination des tâches :

  • ResourceManager : remplace le JobTracker et gère uniquement les ressources globales du cluster, notamment l'allocation des containers (ensemble de ressources CPU, mémoire, réseau, etc.).
  • ApplicationMaster : une instance par application (job ou groupe de jobs) qui gère la coordination et l'exécution des tâches spécifiques à cette application.
  • NodeManager : agent présent sur chaque nœud, responsable de la gestion locale des containers et de la communication avec le ResourceManager.

Cette architecture décentralisée permet une meilleure scalabilité, avec un ResourceManager unique supervisant plusieurs ApplicationMasters, eux-mêmes gérant plusieurs jobs. L'ApplicationMaster négocie les ressources nécessaires auprès du ResourceManager et supervise l'exécution des tâches sur les containers alloués.

Le déroulement d'une tâche sous YARN est le suivant :

  1. Le client soumet un job au ResourceManager, incluant le fichier .jar et la classe driver.
  2. Le ResourceManager alloue un container pour l'ApplicationMaster et lance la classe driver.
  3. L'ApplicationMaster confirme son démarrage et demande des containers pour exécuter les tâches MAP et REDUCE.
  4. Les tâches s'exécutent sur les containers, communiquant leur état à l'ApplicationMaster.
  5. Le client peut interagir directement avec l'ApplicationMaster pour suivre la progression.
  6. À la fin du job, l'ApplicationMaster s'arrête et libère les ressources.

Améliorations de Hadoop 2 et écosystème

Hadoop 2 apporte également des améliorations à HDFS, notamment la suppression du point de défaillance unique (SPOF) du NameNode grâce à un mode actif/attente et la fédération permettant plusieurs NameNodes sur un même cluster.

Cette évolution architecturale découple Hadoop de MapReduce, ouvrant la voie à d'autres frameworks compatibles avec Hadoop et HDFS, et transformant Hadoop en un véritable système d'exploitation pour la donnée.

L'écosystème Hadoop comprend plusieurs outils complémentaires :

  • Pig : langage de haut niveau pour l'analyse de gros volumes de données, destiné aux développeurs habitués aux scripts.
  • Mahout : système d'apprentissage automatique et d'analyse de données implémentant des algorithmes de classification et de clustering.
  • Hive : entrepôt de données avec un langage proche de SQL pour requêter et analyser les données stockées dans Hadoop.
  • HBase : base de données NoSQL orientée colonnes, distribuée et scalable, offrant un accès temps réel.
  • Flume : service distribué pour collecter, agréger et déplacer efficacement de grandes quantités de logs.
  • Zookeeper : service centralisé pour la gestion de la configuration, la synchronisation distribuée et la coordination de services.
  • Sqoop : outil pour transférer efficacement des données entre Hadoop et des bases relationnelles.
  • Oozie : système de gestion de workflows pour coordonner et automatiser les traitements MapReduce et autres tâches.

Points clés à retenir

  • MapReduce est un paradigme de programmation distribué basé sur deux opérations : MAP (transformation en couples clef/valeur) et REDUCE (agrégation des valeurs par clef).
  • Le traitement MapReduce est parallélisable et s'exécute efficacement sur des clusters via Hadoop.
  • Hadoop utilise un JobTracker unique pour gérer les tâches et plusieurs TaskTrackers pour exécuter les opérations sur les nœuds du cluster.
  • YARN, évolution de Hadoop, sépare la gestion des ressources (ResourceManager) de la coordination des jobs (ApplicationMaster), améliorant la scalabilité.
  • Les exemples concrets incluent le comptage de mots dans un texte, les statistiques web et le calcul d'amis en commun dans un réseau social.
  • Hadoop 2 améliore la résilience de HDFS et ouvre l'écosystème à de nombreux outils complémentaires pour le Big Data.

Partager

Commentaires

Aucun commentaire pour le moment. Posez la première question.

Les commentaires sont relus avant publication. Votre e-mail n'est jamais affiché.

← Toutes les révisions