Tous les produits
Search
Centre de documentation

MaxCompute:MultipleInOut example

Dernière mise à jour :Aug 10, 2026

Cet exemple détaille l'exécution d'une tâche MapReduce MultipleInOut sur MaxCompute. La tâche lit les données depuis deux tables d'entrée et dirige la sortie vers deux tables de sortie distinctes — l'une non partitionnée et l'autre partitionnée — en utilisant des libellés de sortie pour contrôler le rédacteur affecté à chaque destination.

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Effectué la configuration de l'environnement décrite dans Prise en main

  • Le fichier mapreduce-examples.jar, stocké dans le répertoire bin\data\resources de votre installation du client MaxCompute

Configurer les tables de test et les ressources

Créer les tables

Exécutez les instructions suivantes dans votre client MaxCompute pour créer deux tables d'entrée et deux tables de sortie. La table mr_multiinout_out2 est partitionnée ; ajoutez-y deux partitions avant d'exécuter la tâche.

CREATE TABLE wc_in1(key STRING, value STRING);
CREATE TABLE wc_in2(key STRING, value STRING);
CREATE TABLE mr_multiinout_out1 (key STRING, cnt BIGINT);
CREATE TABLE mr_multiinout_out2 (key STRING, cnt BIGINT)  PARTITIONED BY (a string, b string);
ALTER TABLE mr_multiinout_out2 ADD PARTITION (a='1', b='1');
ALTER TABLE mr_multiinout_out2 ADD PARTITION (a='2', b='2');

Ajouter la ressource JAR

-- When adding the JAR package for the first time, you can ignore the -f flag.
add jar data\resources\mapreduce-examples.jar -f;

Charger les données de test

Utilisez Tunnel pour charger les fichiers data1.txt et data2.txt depuis le répertoire bin de votre client MaxCompute vers les deux tables d'entrée.

tunnel upload data1.txt wc_in1;
tunnel upload data2.txt wc_in2;

Après le chargement, les tables contiennent les données suivantes :

wc_in1 :

hello,odps

wc_in2 :

hello,world

Exécuter la tâche

Exécutez la commande suivante dans votre client MaxCompute :

jar -resources mapreduce-examples.jar -classpath data\resources\mapreduce-examples.jar
com.aliyun.odps.mapred.open.example.MultipleInOut wc_in1,wc_in2 mr_multiinout_out1,mr_multiinout_out2|a=1/b=1|out1,mr_multiinout_out2|a=2/b=2|out2;

La commande accepte deux arguments positionnels :

Argument

Format

Exemple

Entrées

Noms de table séparés par des virgules

wc_in1,wc_in2

Sorties

Entrées séparées par des virgules, chacune au format table_name|partition_spec|label

mr_multiinout_out2|a=1/b=1|out1

Pour les sorties, les paramètres partition_spec et label sont facultatifs. L'omission du paramètre label dirige la sortie vers la destination par défaut (sans libellé).

Résultats attendus

Le réducteur achemine chaque clé vers une table de sortie spécifique selon le reste de la division du nombre total d'occurrences de la clé par 3 :

Condition

Destination

Table de sortie

count % 3 == 0

Sortie par défaut

mr_multiinout_out1

count % 3 == 1

Libellé out1

mr_multiinout_out2, partition a=1/b=1

count % 3 == 2

Libellé out2

mr_multiinout_out2, partition a=2/b=2

La méthode cleanup() écrit également un enregistrement fixe dans chaque sortie après le traitement de toutes les clés : ("default", 1L), ("out1", 1L) et ("out2", 1L).

mr_multiinout_out1 — reçoit l'enregistrement ("default", 1L) provenant de cleanup(). Aucune clé des données de test ne satisfait la condition count % 3 == 0 ; cette table ne contient donc qu'une seule ligne :

+------------+------------+
| key        | cnt        |
+------------+------------+
| default    | 1          |
+------------+------------+

mr_multiinout_out2 — reçoit les clés dont count % 3 == 1 dans la partition a=1/b=1 et celles dont count % 3 == 2 dans la partition a=2/b=2, ainsi que les enregistrements fixes issus de cleanup() pour chaque partition :

+--------+------------+---+---+
| key    | cnt        | a | b |
+--------+------------+---+---+
| odps   | 1          | 1 | 1 |
| world  | 1          | 1 | 1 |
| out1   | 1          | 1 | 1 |
| hello  | 2          | 2 | 2 |
| out2   | 1          | 2 | 2 |
+--------+------------+---+---+

Exemple de code

Pour la configuration des dépendances Project Object Model (POM), consultez la section Précautions.

L'exemple comprend trois parties principales : le mappeur, le réducteur avec routage multi-sortie et la méthode main qui analyse les arguments des tables d'entrée et de sortie.

Mappeur : analyser les enregistrements d'entrée

Le composant TokenizerMapper émet chaque valeur de champ sous forme de clé avec un compteur de 1.

public static class TokenizerMapper extends MapperBase {
    Record word;
    Record one;

    @Override
    public void setup(TaskContext context) throws IOException {
        word = context.createMapOutputKeyRecord();
        one = context.createMapOutputValueRecord();
        one.set(new Object[] { 1L });
    }

    @Override
    public void map(long recordNum, Record record, TaskContext context)
        throws IOException {
        for (int i = 0; i < record.getColumnCount(); i++) {
            word.set(new Object[] { record.get(i).toString() });
            context.write(word, one);
        }
    }
}

Réducteur : router la sortie par libellé

Le composant SumReducer crée un enregistrement de sortie pour chaque destination libellée dans la méthode setup(). Dans la méthode reduce(), il dirige chaque clé vers une destination selon la valeur de count % 3. La méthode cleanup() écrit un enregistrement fixe dans chaque sortie après le traitement de toutes les clés.

Les API clés pour le routage multi-sortie sont les suivantes :

API

Objectif

context.createOutputRecord()

Créer un enregistrement pour la sortie par défaut (sans libellé)

context.createOutputRecord("out1")

Créer un enregistrement pour la sortie portant le libellé "out1"

context.write(result)

Écrire dans la sortie par défaut

context.write(result1, "out1")

Écrire dans la sortie portant le libellé "out1"

public static class SumReducer extends ReducerBase {
    private Record result;
    private Record result1;
    private Record result2;

    @Override
    public void setup(TaskContext context) throws IOException {
        /** Create a record for each output and add labels to distinguish outputs. */
        result = context.createOutputRecord();          // default output (no label)
        result1 = context.createOutputRecord("out1");  // output labeled "out1"
        result2 = context.createOutputRecord("out2");  // output labeled "out2"
    }

    @Override
    public void reduce(Record key, Iterator<Record> values, TaskContext context)
        throws IOException {
        long count = 0;
        while (values.hasNext()) {
            Record val = values.next();
            count += (Long) val.get(0);
        }
        long mod = count % 3;
        if (mod == 0) {
            result.set(0, key.get(0));
            result.set(1, count);
            /** If you do not specify a label, the default output is used. */
            context.write(result);
        } else if (mod == 1) {
            result1.set(0, key.get(0));
            result1.set(1, count);
            context.write(result1, "out1");
        } else {
            result2.set(0, key.get(0));
            result2.set(1, count);
            context.write(result2, "out2");
        }
    }

    @Override
    public void cleanup(TaskContext context) throws IOException {
        Record result = context.createOutputRecord();
        result.set(0, "default");
        result.set(1, 1L);
        context.write(result);

        Record result1 = context.createOutputRecord("out1");
        result1.set(0, "out1");
        result1.set(1, 1L);
        context.write(result1, "out1");

        Record result2 = context.createOutputRecord("out2");
        result2.set(0, "out2");
        result2.set(1, 1L);
        context.write(result2, "out2");
    }
}

Méthode principale : configurer les entrées et les sorties

La méthode main analyse les chaînes de caractères des tables d'entrée et de sortie à partir des arguments de ligne de commande, puis soumet la tâche. Les entrées utilisent le format table_name ou table_name|partition_spec. Les sorties ajoutent un libellé facultatif : table_name|partition_spec|label.

/** Convert partition strings such as "ds=1/pt=2" to MAP. */
public static LinkedHashMap<String, String> convertPartSpecToMap(String partSpec) {
    LinkedHashMap<String, String> map = new LinkedHashMap<String, String>();
    if (partSpec != null && !partSpec.trim().isEmpty()) {
        String[] parts = partSpec.split("/");
        for (String part : parts) {
            String[] ss = part.split("=");
            if (ss.length != 2) {
                throw new RuntimeException(
                    "ODPS-0730001: error part spec format: " + partSpec);
            }
            map.put(ss[0], ss[1]);
        }
    }
    return map;
}

public static void main(String[] args) throws Exception {
    String[] inputs = null;
    String[] outputs = null;
    if (args.length == 2) {
        inputs = args[0].split(",");
        outputs = args[1].split(",");
    } else {
        System.err.println("MultipleInOut in... out...");
        System.exit(1);
    }

    JobConf job = new JobConf();
    job.setMapperClass(TokenizerMapper.class);
    job.setReducerClass(SumReducer.class);
    job.setMapOutputKeySchema(SchemaUtils.fromString("word:string"));
    job.setMapOutputValueSchema(SchemaUtils.fromString("count:bigint"));

    /** Parse input table strings. */
    for (String in : inputs) {
        String[] ss = in.split("\\|");
        if (ss.length == 1) {
            InputUtils.addTable(TableInfo.builder().tableName(ss[0]).build(), job);
        } else if (ss.length == 2) {
            LinkedHashMap<String, String> map = convertPartSpecToMap(ss[1]);
            InputUtils.addTable(
                TableInfo.builder().tableName(ss[0]).partSpec(map).build(), job);
        } else {
            System.err.println("Style of input: " + in + " is not right");
            System.exit(1);
        }
    }

    /** Parse output table strings. */
    for (String out : outputs) {
        String[] ss = out.split("\\|");
        if (ss.length == 1) {
            OutputUtils.addTable(TableInfo.builder().tableName(ss[0]).build(), job);
        } else if (ss.length == 2) {
            LinkedHashMap<String, String> map = convertPartSpecToMap(ss[1]);
            OutputUtils.addTable(
                TableInfo.builder().tableName(ss[0]).partSpec(map).build(), job);
        } else if (ss.length == 3) {
            if (ss[1].isEmpty()) {
                LinkedHashMap<String, String> map = convertPartSpecToMap(ss[2]);
                OutputUtils.addTable(
                    TableInfo.builder().tableName(ss[0]).partSpec(map).build(), job);
            } else {
                LinkedHashMap<String, String> map = convertPartSpecToMap(ss[1]);
                OutputUtils.addTable(
                    TableInfo.builder().tableName(ss[0]).partSpec(map)
                        .label(ss[2]).build(), job);
            }
        } else {
            System.err.println("Style of output: " + out + " is not right");
            System.exit(1);
        }
    }

    JobClient.runJob(job);
}

Imports complets pour cet exemple :

package com.aliyun.odps.mapred.open.example;
import java.io.IOException;
import java.util.Iterator;
import java.util.LinkedHashMap;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.data.TableInfo;
import com.aliyun.odps.mapred.JobClient;
import com.aliyun.odps.mapred.MapperBase;
import com.aliyun.odps.mapred.ReducerBase;
import com.aliyun.odps.mapred.TaskContext;
import com.aliyun.odps.mapred.conf.JobConf;
import com.aliyun.odps.mapred.utils.InputUtils;
import com.aliyun.odps.mapred.utils.OutputUtils;
import com.aliyun.odps.mapred.utils.SchemaUtils;