MapReduce

Cette conférence présente le paradigme MapReduce, un modèle de programmation essentiel pour le traitement distribué de grandes quantités de données. Elle s'inscrit dans un cours sur le Big Data, abordant à la fois les principes fondamentaux de MapReduce, des exemples concrets d'application, ainsi que l'architecture Hadoop et son évolution vers YARN (MapReduce 2).

D'après le document MapReduce

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

MapReduce

Document source

MapReduce

Big Data · PDF · 88 pages · 2016

Afficher l'aperçu du document

Consulter le document original →

Cette conférence présente le paradigme MapReduce, un modèle de programmation essentiel pour le traitement distribué de grandes quantités de données. Elle s'inscrit dans un cours sur le Big Data, abordant à la fois les principes fondamentaux de MapReduce, des exemples concrets d'application, ainsi que l'architecture Hadoop et son évolution vers YARN (MapReduce 2).

Introduction à MapReduce

MapReduce est un paradigme de programmation conçu pour traiter efficacement de très grands ensembles de données en les divisant en sous-problèmes plus petits, exécutés en parallèle sur un cluster de machines. Ce modèle s'appuie sur la stratégie algorithmique du "divide and conquer" (diviser pour régner).

Le principe consiste à appliquer deux opérations distinctes :

  • MAP : transforme les données d'entrée en une série de couples (clef, valeur). Cette étape est parallélisable car chaque fragment de données peut être traité indépendamment sur une machine différente.
  • REDUCE : agrège les valeurs associées à chaque clef distincte produite par l'étape MAP, en appliquant un traitement spécifique pour obtenir un résultat final par clef.

Le traitement MapReduce comprend quatre étapes principales :

  1. Découpage (split) des données d'entrée en fragments.
  2. Application de la fonction MAP sur chaque fragment pour générer des couples (clef, valeur).
  3. Regroupement (shuffle) des couples par clef.
  4. Application de la fonction REDUCE sur chaque groupe pour produire la sortie finale.

Cette organisation permet une parallélisation efficace, rendant possible le traitement de 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 déterminer les mots les plus fréquents. Les étapes sont les suivantes :

  • Les données d'entrée sont le contenu du texte, préalablement nettoyé (suppression de la ponctuation, conversion en minuscules, suppression des accents).
  • Le découpage s'effectue ligne par ligne, chaque ligne devenant un fragment indépendant.
  • La clef choisie pour MAP est le mot lui-même.
  • L'opération MAP génère pour chaque mot rencontré le couple (mot ; 1), indiquant une occurrence unique.

Par exemple, pour la ligne "celui qui croyait au ciel", MAP produit :

(celui;1) (qui;1) (croyait;1) (au;1) (ciel;1)

Après l'étape MAP sur tous les fragments, Hadoop effectue automatiquement le regroupement (shuffle) des couples par clef, formant des groupes comme :

(qui;1) (qui;1) (qui;1) (qui;1)

L'opération REDUCE consiste alors à sommer les valeurs associées à chaque mot :

TOTAL = 0
POUR chaque occurrence dans le groupe:
    TOTAL = TOTAL + 1
RENVOYER TOTAL

Le résultat final indique les fréquences des mots, par exemple :

qui : 4
celui : 2
croyait : 2
fou : 2
au : 1
ciel : 1
ny : 1
pas : 1
fait : 1
[…]

Ce modèle, bien que simple dans cet exemple, est applicable à des corpus beaucoup plus volumineux, comme l'ensemble des textes d'une bibliothèque, avec un traitement distribué efficace.

Autres exemples d'application de MapReduce

Statistiques web

Un autre cas d'usage consiste à compter le nombre de visiteurs par page d'un site Internet. Les données d'entrée sont les fichiers de logs contenant les URL des pages visitées. La clef MAP est l'URL, et les opérations MAP et REDUCE sont similaires à l'exemple précédent, permettant de calculer le nombre de vues par page.

Calcul des amis en commun dans un réseau social

Considérons un réseau social avec des millions d'utilisateurs, où chaque utilisateur a une liste d'amis. L'objectif est d'afficher, lorsqu'un utilisateur visite la page d'un autre, le nombre d'amis qu'ils ont en commun. Pour éviter des requêtes SQL coûteuses en temps réel, on utilise MapReduce pour pré-calculer ces données.

Les données d'entrée sont sous la forme : Utilisateur => Liste d'amis.

La clef choisie est la concaténation ordonnée alphabétiquement de deux utilisateurs, par exemple "A-B" pour les utilisateurs A et B. L'opération MAP génère pour chaque utilisateur toutes les paires possibles avec ses amis, associant à chaque clef la liste d'amis de l'utilisateur :

Pseudo-code MAP:
UTILISATEUR = première partie de la ligne
POUR AMI dans liste d'amis:
    SI UTILISATEUR < AMI:
        CLEF = UTILISATEUR + "-" + AMI
    SINON:
        CLEF = AMI + "-" + UTILISATEUR
    GENERER COUPLE (CLEF; liste d'amis)

Par exemple, pour la ligne "A => B, C, D", MAP produit :

("A-B"; "B C D")
("A-C"; "B C D")
("A-D"; "B C D")

Après regroupement, chaque clef possède deux listes d'amis, une pour chaque utilisateur. L'opération REDUCE consiste à calculer l'intersection de ces deux listes, c'est-à-dire les amis communs :

LISTE_AMIS_COMMUNS = []
SI longueur(VALEURS) != 2:
    RENVOYER ERREUR
SINON:
    POUR AMI dans VALEURS[0]:
        SI AMI dans VALEURS[1]:
            LISTE_AMIS_COMMUNS += AMI
RENVOYER LISTE_AMIS_COMMUNS

Le résultat donne, par exemple :

"A-B" : "C, D"
"A-C" : "B, D"
"B-C" : "A, D, E"
[…]

Ce traitement, parallélisable sur un cluster Hadoop, permet de gérer efficacement des millions d'utilisateurs.

Architecture Hadoop pour MapReduce

Hadoop est une plateforme open source qui implémente MapReduce et le système de fichiers distribué HDFS. L'architecture Hadoop repose sur deux types principaux de serveurs :

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

Le JobTracker attribue les tâches aux TaskTracker en fonction de la proximité des données (localité), ce qui optimise les performances. Il surveille également la progression des tâches via des messages "heartbeat" et peut relancer les tâches en cas d'échec.

Chaque TaskTracker dispose d'un nombre configurable de "slots" d'exécution correspondant au nombre de tâches pouvant être traitées simultanément sur la machine.

Le JobTracker et le NameNode sont généralement déployés sur une machine maîtresse dédiée, tandis que les autres machines hébergent les TaskTracker et DataNodes.

Évolution vers YARN (MapReduce 2)

YARN (Yet Another Resource Negotiator) est une évolution majeure de l'architecture MapReduce, introduite dans Hadoop 2. Elle répond aux limites de scalabilité et de gestion des ressources du JobTracker unique.

Les responsabilités du JobTracker sont réparties entre :

  • ResourceManager : gère globalement les ressources du cluster et l'allocation des conteneurs (containers) qui regroupent CPU, mémoire, disque et réseau.
  • ApplicationMaster : une instance par application (job) qui gère l'exécution des tâches, négocie les ressources auprès du ResourceManager et supervise la progression.
  • NodeManager : agent sur chaque nœud qui supervise l'utilisation des ressources et exécute les conteneurs.

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

  1. Le client soumet une application (archive .jar et classe driver) au ResourceManager.
  2. Le ResourceManager alloue un container pour l'ApplicationMaster et lance la classe driver.
  3. L'ApplicationMaster confirme son démarrage et demande au ResourceManager des containers pour exécuter les fragments de données.
  4. L'ApplicationMaster lance les tâches MAP et REDUCE sur les containers alloués, communiquant directement avec eux.
  5. Les NodeManagers rapportent régulièrement l'état des ressources au ResourceManager.
  6. Le client peut interagir directement avec l'ApplicationMaster pour suivre la progression.
  7. À la fin du job, l'ApplicationMaster s'arrête et libère les ressources.

Cette architecture découplée améliore la scalabilité, la flexibilité et la gestion des ressources, permettant d'exécuter plusieurs applications simultanément sur un même cluster.

Améliorations de Hadoop 2 et écosystème

Hadoop 2 introduit également des améliorations dans HDFS, notamment :

  • Suppression du SPOF (Single Point Of Failure) du NameNode grâce à une architecture active/attente.
  • Fédération HDFS permettant plusieurs NameNodes pour gérer différents espaces de noms sur un même cluster.

Hadoop devient ainsi un véritable système d'exploitation pour les données, supportant plusieurs frameworks et langages.

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

  • Pig : langage de haut niveau pour l'analyse de données volumineuses, destiné aux développeurs familiers avec les scripts.
  • Mahout : bibliothèque d'apprentissage automatique et de data mining, avec des algorithmes de classification et de clustering.
  • Hive : entrepôt de données avec un langage de requête proche de SQL pour l'analyse et l'agrégation de données Hadoop.
  • HBase : base de données NoSQL orientée colonnes, distribuée et scalable, pour 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 des 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 des traitements MapReduce et autres tâches.

Points clés à retenir

  • MapReduce est un modèle de programmation distribué basé sur deux opérations : MAP (transformation en couples clef/valeur) et REDUCE (agrégation par clef).
  • Le traitement MapReduce est parallélisable et permet de gérer efficacement de très grands volumes de données.
  • Hadoop implémente MapReduce avec une architecture maître-esclave : JobTracker (maître) et TaskTracker (esclaves).
  • YARN (MapReduce 2) sépare la gestion des ressources (ResourceManager) de la gestion des applications (ApplicationMaster), améliorant la scalabilité et la flexibilité.
  • Hadoop 2 améliore HDFS avec une architecture tolérante aux pannes et la fédération des NameNodes.
  • L'écosystème Hadoop comprend de nombreux outils facilitant l'analyse, le stockage, la gestion et le transfert des données massives.

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