Tous les produits
Search
Centre de documentation

MaxCompute:Vue d'ensemble

Dernière mise à jour :Aug 10, 2026

Cette rubrique décrit les classes et méthodes principales de MapReduce dans le SDK Java MaxCompute (odps-sdk-mapred).

Ajout de la dépendance du SDK

Recherchez odps-sdk-mapred dans le référentiel Maven pour identifier la dernière version. Ajoutez la dépendance suivante à votre projet :

<dependency>
    <groupId>com.aliyun.odps</groupId>
    <artifactId>odps-sdk-mapred</artifactId>
    <version>0.40.10-public</version>
</dependency>

Types de données

MaxCompute MapReduce prend en charge les types de données suivants. Le tableau ci-dessous indique leur correspondance avec les types Java.

Type MaxCompute

Type Java

BIGINT

Long

STRING

String

DOUBLE

Double

BOOLEAN

Boolean

DATETIME

Date

DECIMAL

BigDecimal

Présentation des classes

Classe

Description

MapperBase

Classe de base pour les mappeurs définis par l'utilisateur. Elle convertit les enregistrements des tables d'entrée en paires clé-valeur et les transmet au réducteur. Les tâches qui ignorent l'étape de réduction et écrivent directement les résultats sont appelées tâches MapOnly.

ReducerBase

Classe de base pour les réducteurs définis par l'utilisateur. Elle agrège l'ensemble des valeurs associées à chaque clé.

TaskContext

Fournit le contexte d'exécution de la tâche. Cet objet est transmis en tant que paramètre d'entrée aux méthodes de cycle de vie de MapperBase et ReducerBase.

JobClient

Soumet et gère les tâches. Il prend en charge la soumission bloquante (synchrone) et non bloquante (asynchrone).

RunningJob

Représente une instance de tâche en cours d'exécution. Utilisez cet objet pour suivre l'état, attendre la fin de l'exécution et récupérer les compteurs.

JobConf

Contient la configuration d'une tâche MapReduce. Définissez un objet JobConf dans votre fonction principale, puis transmettez-le à JobClient pour soumettre la tâche.

MapperBase

Le framework appelle les méthodes du cycle de vie du mappeur dans l'ordre suivant : setup une fois au début, map une fois par enregistrement d'entrée, et cleanup une fois à la fin.

public class WordCountMapper extends MapperBase {
    @Override
    public void setup(TaskContext context) throws IOException {
        // Initialize resources (for example, load lookup tables)
    }

    @Override
    public void map(long key, Record record, TaskContext context) throws IOException {
        // Process each input record and emit key-value pairs.
        // MaxCompute uses (long key, Record record), not (KEYIN key, VALUEIN value).
        Record mapKey = context.createMapOutputKeyRecord();
        Record mapValue = context.createMapOutputValueRecord();
        // Populate mapKey and mapValue, then emit:
        context.write(mapKey, mapValue);
    }

    @Override
    public void cleanup(TaskContext context) throws IOException {
        // Release resources or emit final accumulated results
    }
}

Méthode

Description

void setup(TaskContext context)

Appelée une fois avant le premier appel à map. Utilisez cette méthode pour initialiser l'état partagé ou charger des ressources.

void map(long key, Record record, TaskContext context)

Appelée une fois par enregistrement d'entrée. Émettez des paires clé-valeur avec context.write(key, value) ou écrivez directement dans une table de sortie.

void cleanup(TaskContext context)

Appelée une fois après le dernier appel à map. Utilisez cette méthode pour libérer des ressources ou vider la sortie mise en tampon.

ReducerBase

Le framework appelle les méthodes du cycle de vie du réducteur dans l'ordre suivant : setup une fois au début, reduce une fois par groupe de clés uniques, et cleanup une fois à la fin.

public class WordCountReducer extends ReducerBase {
    @Override
    public void setup(TaskContext context) throws IOException {
        // Initialize resources
    }

    @Override
    public void reduce(Record key, Iterator<Record> values, TaskContext context) throws IOException {
        // Aggregate all values for this key
        long count = 0;
        while (values.hasNext()) {
            Record val = values.next();
            count += val.getBigint(0);
        }
        Record output = context.createOutputRecord();
        output.set(0, key.getString(0));
        output.set(1, count);
        context.write(output);
    }

    @Override
    public void cleanup(TaskContext context) throws IOException {
        // Release resources
    }
}

Méthode

Description

void setup(TaskContext context)

Appelée une fois avant le premier appel à reduce. Utilisez cette méthode pour initialiser l'état partagé ou charger des ressources.

void reduce(Record key, Iterator<Record> values, TaskContext context)

Appelée une fois par groupe de clés uniques. Toutes les valeurs associées à la clé sont transmises via values.

void cleanup(TaskContext context)

Appelée une fois après le dernier appel à reduce. Utilisez cette méthode pour libérer des ressources ou vider la sortie mise en tampon.

TaskContext

L'objet TaskContext est transmis à chaque méthode du cycle de vie et permet d'accéder aux tables de sortie, aux ressources et aux utilitaires du framework.

Méthode

Description

TableInfo[] getOutputTableInfo()

Renvoie les informations relatives aux tables de sortie.

Record createOutputRecord()

Crée un enregistrement pour la table de sortie par défaut.

Record createOutputRecord(String label)

Crée un enregistrement pour la table de sortie identifiée par label.

Record createMapOutputKeyRecord()

Crée un enregistrement pour la clé de sortie de l'étape map.

Record createMapOutputValueRecord()

Crée un enregistrement pour la valeur de sortie de l'étape map.

void write(Record record)

Écrit un enregistrement dans la table de sortie par défaut. Cette méthode peut être appelée plusieurs fois pendant l'étape de réduction.

void write(Record record, String label)

Écrit un enregistrement dans la table de sortie identifiée par label. Cette méthode peut être appelée plusieurs fois pendant l'étape de réduction.

void write(Record key, Record value)

Émet une paire clé-valeur pendant l'étape map. Cette méthode peut être appelée plusieurs fois.

BufferedInputStream readResourceFileAsStream(String resourceName)

Lit une ressource de fichier par son nom.

Iterator<Record> readResourceTable(String resourceName)

Lit une ressource de table par son nom.

Counter getCounter(Enum<?> name)

Renvoie le compteur portant le nom spécifié.

Counter getCounter(String group, String name)

Renvoie le compteur portant le nom spécifié dans le groupe indiqué.

void progress()

Envoie un signal de maintien actif (heartbeat) au framework MapReduce pour éviter le dépassement du délai d'inactivité du worker.

Délai d'inactivité du worker

Le délai d'inactivité par défaut du worker est de 10 minutes et ne peut pas être modifié. Si un worker n'appelle pas progress() dans un délai de 10 minutes, le framework met fin au processus du worker et la tâche map ou reduce échoue. Appelez progress() périodiquement dans les tâches de longue durée pour maintenir les workers actifs. Cette méthode envoie uniquement un signal de maintien actif ; elle ne signale pas la progression de la tâche.

JobConf

L'objet JobConf contient toute la configuration d'une tâche MapReduce, y compris les classes de mappeur et de réducteur, les schémas clé-valeur et les déclarations de ressources.

Méthode

Description

void setMapperClass(Class<? extends Mapper> theClass)

Définit la classe du mappeur pour la tâche.

void setReducerClass(Class<? extends Reducer> theClass)

Définit la classe du réducteur pour la tâche.

void setCombinerClass(Class<? extends Reducer> theClass)

Définit une classe de combiner. Un combiner pré-agrège les enregistrements partageant la même clé lors de l'étape map, ce qui réduit le volume de données transféré vers les réducteurs.

void setMapOutputKeySchema(Column[] schema)

Définit le schéma des clés transmises du mappeur au réducteur.

void setMapOutputValueSchema(Column[] schema)

Définit le schéma des valeurs transmises du mappeur au réducteur.

void setOutputKeySortColumns(String[] cols)

Définit les colonnes utilisées pour trier les clés avant leur transmission aux réducteurs.

void setOutputGroupingColumns(String[] cols)

Définit les colonnes utilisées pour regrouper les clés. Les colonnes de regroupement doivent constituer un sous-ensemble des colonnes de tri.

void setPartitionColumns(String[] cols)

Définit les colonnes de clé de partition. Par défaut, toutes les colonnes de clé sont utilisées.

void setResources(String resourceNames)

Déclare les ressources accessibles aux mappeurs et aux réducteurs. Un mappeur ou un réducteur ne peut lire que les ressources déclarées ici.

void setSplitSize(long size)

Définit la taille du fragment d'entrée en Mo. Valeur par défaut : 256 Mo.

void setNumReduceTasks(int n)

Définit le nombre de tâches de réduction. Valeur par défaut : un quart du nombre de tâches map.

void setMemoryForMapTask(int mem)

Définit la mémoire allouée par worker map en Mo. Valeur par défaut : 2048 Mo.

void setMemoryForReduceTask(int mem)

Définit la mémoire allouée par worker reduce en Mo. Valeur par défaut : 2048 Mo.

Distribution et regroupement des enregistrements lors des étapes map et reduce

Lors de l'étape map, le framework calcule un hachage de chaque enregistrement de sortie en fonction des colonnes de clé de partition afin de déterminer quel réducteur doit le recevoir. Les enregistrements sont triés selon les colonnes de tri avant d'être envoyés.

Lors de l'étape reduce, les enregistrements sont regroupés par les colonnes de regroupement. Tous les enregistrements partageant la même clé de regroupement sont transmis ensemble à un seul appel de reduce().

Les colonnes de regroupement sont sélectionnées parmi les colonnes de tri. Les colonnes de tri et les colonnes de clé de partition doivent exister dans les clés.

JobClient

Méthode

Description

static RunningJob runJob(JobConf job)

Soumet une tâche en mode bloquant (synchrone). L'appel bloque jusqu'à la fin de la tâche et renvoie un objet RunningJob représentant la tâche terminée.

static RunningJob submitJob(JobConf job)

Soumet une tâche en mode non bloquant (asynchrone). L'appel renvoie immédiatement un objet RunningJob que vous pouvez utiliser pour interroger l'état ou attendre la fin de l'exécution.

RunningJob

Méthode

Description

String getInstanceID()

Renvoie l'ID de l'instance de tâche. Utilisez cet ID pour consulter les journaux opérationnels et gérer la tâche.

boolean isComplete()

Renvoie true si la tâche est terminée.

boolean isSuccessful()

Renvoie true si la tâche s'est terminée avec succès.

void waitForCompletion()

Attend la fin d'une instance de tâche. Cette méthode est utilisée pour les tâches soumises en mode synchrone.

JobStatus getJobStatus()

Renvoie l'état actuel de l'instance de tâche.

void killJob()

Met fin à la tâche en cours d'exécution.

Counters getCounters()

Renvoie toutes les données des compteurs pour la tâche.

InputUtils

Méthode

Description

static void addTable(TableInfo table, JobConf conf)

Ajoute une seule table d'entrée à la tâche. Des appels multiples à cette méthode ajoutent chaque table à la file d'attente d'entrée.

static void setTables(TableInfo[] tables, JobConf conf)

Définit plusieurs tables d'entrée en une seule fois.

OutputUtils

Méthode

Description

static void addTable(TableInfo table, JobConf conf)

Ajoute une seule table de sortie à la tâche. Des appels multiples à cette méthode ajoutent chaque table à la file d'attente de sortie.

static void setTables(TableInfo[] tables, JobConf conf)

Définit plusieurs tables de sortie en une seule fois.

Pipeline (modèle MapReduce étendu)

L'objet Pipeline constitue le point d'entrée du modèle MapReduce étendu, qui permet d'enchaîner plusieurs mappeurs et réducteurs au sein d'une seule tâche. Utilisez Pipeline.builder() pour construire le pipeline.

Méthodes du générateur

public Builder addMapper(Class<? extends Mapper> mapper)
public Builder addMapper(Class<? extends Mapper> mapper,
       Column[] keySchema, Column[] valueSchema, String[] sortCols,
       SortOrder[] order, String[] partCols,
       Class<? extends Partitioner> theClass, String[] groupCols)
public Builder addReducer(Class<? extends Reducer> reducer)
public Builder addReducer(Class<? extends Reducer> reducer,
       Column[] keySchema, Column[] valueSchema, String[] sortCols,
       SortOrder[] order, String[] partCols,
       Class<? extends Partitioner> theClass, String[] groupCols)
public Builder setOutputKeySchema(Column[] keySchema)
public Builder setOutputValueSchema(Column[] valueSchema)
public Builder setOutputKeySortColumns(String[] sortCols)
public Builder setOutputKeySortOrder(SortOrder[] order)
public Builder setPartitionColumns(String[] partCols)
public Builder setPartitionerClass(Class<? extends Partitioner> theClass)
public Builder setOutputGroupingColumns(String[] cols)

Exemple

L'exemple suivant enchaîne un mappeur et deux réducteurs (TokenizerMapperSumReducerIdentityReducer) :

Job job = new Job();
Pipeline pipeline = Pipeline.builder()
    .addMapper(TokenizerMapper.class)
    .setOutputKeySchema(
        new Column[] { new Column("word", OdpsType.STRING) })
    .setOutputValueSchema(
        new Column[] { new Column("count", OdpsType.BIGINT) })
    .addReducer(SumReducer.class)
    .setOutputKeySchema(
        new Column[] { new Column("count", OdpsType.BIGINT) })
    .setOutputValueSchema(
        new Column[] { new Column("word", OdpsType.STRING),
        new Column("count", OdpsType.BIGINT) })
    .addReducer(IdentityReducer.class).createPipeline();

job.setPipeline(pipeline);
job.addInput(...)
job.addOutput(...)
job.submit();
Pour enchaîner un seul mappeur avec un réducteur, utilisez JobConf au lieu de Pipeline . Familiarisez-vous avec l'API MapReduce standard avant d'utiliser le modèle étendu.