Architectures de tri global sous Apache Hadoop : partitionnement intelligent et échantillonnage adaptatif

Principe du tri distribué et limitation du réducteur unique Obtenir une sortie globalement ordonnée dans un cluster Hadoop pose un défi architectural majeur. La solution naïve consiste à affecter un seul réducteur à la tâche de tri. Bien que cette approche garantisse un ordre total, elle annule totalement les bénéfices de l'infrastructure parallèle. Le col de bouteille devient immédiat lorsque le volume de données dépasse la capacité mémoire ou CPU d'un nœud unique, provoquant des goulots d'étranglement sévères et une augmentation exponentielle du temps d'exécution.

Architecture basée sur le partitionnement par intervalles Pour préserver la scalabilité, il est impératif de répartir les données entre plusieurs réducteurs tout en conservant l'ordre global au niveau du cluster. Cette méthode repose sur le partitionnement hiérarchique : les clés sont découpées selon des bornes précalculées. Chaque groupe de valeurs est ensuite envoyé vers un réducteur spécifique qui trie localement son lot. À la fin de l'exécution, chaque fichier de sortie correspondant à un réducteur contient des données ordonnées, et l'ensemble des fichiers respecte un ordre croissant entre eux.

Cependant, une distribution arbitraire des intervalles génère souvent une déséquilibration des charges. Certaines partitions peuvent contenir 80 % du jeu de données tandis que d'autres reçoivent moins de 5 %. Pour optimiser les performances, il faut viser une répartition homogène du nombre d'enregistrements par partition, évitant ainsi qu'une tâche lente ne bloque l'orchestration globale.

Stratégies d'échantillonnage pour déterminer les seuils Définir dynamiquement les bornes sans parcourir l'intégralité du dataset serait extrêmement coûteux. Hadoop propose donc un module dédié : InputSampler. Celui-ci extrait un sous-ensemble représentatif des entrées, calcule statistiquement les séparateurs idéaux, et génère un fichier de configuration lu ultérieurement par TotalOrderPartitioner. Trois algorithmes d'extraction coexistent :

Méthode Principe Paramètres requis Performances
Séquentiel (SplitSampler) Extrait les premières valeurs de chaque division logique cibleTotale, divisionsMax Optimale
Aléatoire (RandomSampler) Parcours complet avec probabilité constante et remplacement progressif taux抽取, cibleTotale, divisionsMax Modérée
Intervalle fixe (IntervalSampler) Collecte périodique selon un ratio temporel ou comptable taux抽取, divisionsMax Intermédiaire

Fonctionnement interne du générateur de partision Le mécanisme central s'appuie sur l'interface Sampler fournie par la bibliothèque mapreduce.lib.partition. Après récupération des fragments, l'orchestrateur trie les extraits, calcule le pas de saute (stepSize), et sélectionne les valeurs pivot équidistantes. Ces dernières sont sérialisées dans un fichier binaire indexé, puis diffusées via le cache distribué du framework. Au moment de l'exécution mapper, TotalOrderPartitioner interroge ce cache, détermine le segment approprié via une recherche dichotomique, et assigne la clé au numéro de tâche cible.

Implémentation refactorisée Voici une version réorganisée de la classe principale, intégrant une configuration modulaire et une gestion explicite des ressources. Les variables et la structure de contrôle ont été adaptées pour améliorer la lisibilité tout en respectant strictement les contrats de l'API MapReduce.

package org.apache.hadoop.examples.globalorder;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.compress.GzipCodec;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.SequenceFileOutputFormat;
import org.apache.hadoop.mapreduce.lib.partition.InputSampler;
import org.apache.hadoop.mapreduce.lib.partition.TotalOrderPartitioner;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;

import java.net.URI;

public class GlobalSortPipeline implements Tool {

    private Configuration systemConf;

    @Override
    public void setConf(Configuration conf) {
        this.systemConf = conf;
    }

    @Override
    public Configuration getConf() {
        return this.systemConf;
    }

    @Override
    public int run(String[] arguments) throws Exception {
        Job executionTask = Job.getInstance(getConf(), getClass().getSimpleName());
        executionTask.setJarByClass(GlobalSortPipeline.class);

        FileInputFormat.addInputPath(executionTask, new Path(arguments[0]));
        FileOutputFormat.setOutputPath(executionTask, new Path(arguments[1]));

        // Configuration des formats et compression
        executionTask.setOutputKeyClass(IntWritable.class);
        executionTask.setOutputValueClass(Text.class);
        executionTask.setInputFormatClass(org.apache.hadoop.mapreduce.lib.input.SequenceFileInputFormat.class);
        executionTask.setOutputFormatClass(SequenceFileOutputFormat.class);

        SequenceFileOutputFormat.setCompressOutput(executionTask, true);
        SequenceFileOutputFormat.setOutputCompressorClass(executionTask, GzipCodec.class);

        // Activation du partitionneur global
        executionTask.setPartitionerClass(TotalOrderPartitioner.class);

        // Définition du nombre de tâches réductrices cibles
        int targetReducers = Integer.parseInt(systemConf.get("mapper.reduce.count", "4"));
        executionTask.setNumReduceTasks(targetReducers);

        // Initialisation de l'heuristic d'échantillonnage
        double samplingProbability = Double.parseDouble(systemConf.get("sample.rate", "0.1"));
        InputSampler.Sampler<IntWritable, Text> distributionModel = 
            new InputSampler.RandomSampler<>(samplingProbability, 10000, 10);

        // Génération et persistance des frontières de partition
        InputSampler.writePartitionFile(executionTask, distributionModel);

        // Rattachement dynamique du fichier de seuils au contexte distribué
        String partitionIndexUri = TotalOrderPartitioner.getPartitionFile(getConf());
        executionTask.addCacheFile(new URI(partitionIndexUri));

        boolean success = executionTask.waitForCompletion(false);
        return success ? 0 : 1;
    }

    public static void main(String[] args) throws Exception {
        int exitStatus = ToolRunner.run(new GlobalSortPipeline(), args);
        System.exit(exitStatus);
    }
}

Consignes d'exécution Lancez la pipeline via le client CLI en spécifiant explicitement le paramètre de réduction. L'indicateur -totalsort active automatiquement la chaîne d'optimisation globale lors de l'emballage de l'archive exécutable :

hadoop jar lib/hadoop-custom-examples.jar \
  -D mapper.reduce.count=4 \
  -D sample.rate=0.15 \
  GlobalSortPipeline input/raw_dataset output/sorted_results -totalsort

La phase d'échantillonnage s'exécute généralement comme un job préliminaire léger, suivi du traitement principle qui exploite les bornes générées. L'architecture garantit une montée en charge linéaire tout en maintenant la cohérence transactionnelle de l'ordre alphabétique ou numérique des enregistrements sortants.

Étiquettes: HadoopMapReduce InputSampler TotalOrderPartitioner DistributedSorting DataPartitioning

Publié le 25 septembre à 13h40