Tous les produits
Search
Centre de documentation

MaxCompute:Exemples de pipeline

Dernière mise à jour :Aug 10, 2026

Cet exemple montre comment compter les fréquences des mots à l'aide d'une tâche Pipeline MapReduce MaxCompute. Le pipeline enchaîne un mappeur (TokenizerMapper) et deux réducteurs (SumReducer et IdentityReducer), illustrant ainsi la manière dont l'API Pipeline relie plusieurs étapes de traitement au sein d'une même tâche.

Prérequis

Avant de commencer, assurez-vous d'avoir :

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

Configurer les ressources de test

1. Préparer le package JAR

Placez mapreduce-examples.jar dans le répertoire bin\data\resources de votre installation du client MaxCompute.

2. Créer des tables

Exécutez les instructions SQL suivantes dans le client MaxCompute :

CREATE TABLE wc_in (key STRING, value STRING);
CREATE TABLE wc_out(key STRING, cnt BIGINT);

3. Ajouter le package JAR en tant que ressource

add jar data\resources\mapreduce-examples.jar -f;
L'option -f écrase une ressource existante portant le même nom. Omettez-la lors du premier ajout du package JAR.

4. Importer les données de test

Utilisez une commande Tunnel pour charger data.txt depuis le répertoire bin du client MaxCompute vers wc_in :

tunnel upload data.txt wc_in;

Le fichier data.txt contient la ligne suivante :

hello,odps

Exécuter le pipeline

Exécutez la commande suivante dans le client MaxCompute :

jar -resources mapreduce-examples.jar -classpath data\resources\mapreduce-examples.jar
com.aliyun.odps.mapred.open.example.WordCountPipeline wc_in wc_out;

Vérifier le résultat

Une fois la tâche terminée, interrogez wc_out pour confirmer la sortie :

SELECT * FROM wc_out;

La sortie attendue est la suivante :

+------------+------------+
| key        | cnt        |
+------------+------------+
| hello      | 1          |
| odps       | 1          |
+------------+------------+

Exemple de code

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

Fonctionnement du pipeline

Le pipeline enchaîne trois étapes. Chaque étape transforme les données comme suit :

wc_in (key STRING, value STRING)TokenizerMapper(word STRING, count BIGINT)SumReducer(word STRING, count BIGINT)IdentityReducerwc_out (key STRING, cnt BIGINT)

Classe

Rôle

TokenizerMapper

Divise chaque enregistrement en jetons et émet (word, 1) pour chaque jeton

SumReducer

Agrège les comptes pour chaque clé de mot

IdentityReducer

Écrit les paires finales (word, count) dans la table de sortie

La conception à deux réducteurs sépare l'agrégation du formatage de la sortie. SumReducer accumule les totaux provenant de toutes les sorties des mappeurs. IdentityReducer mappe ensuite le schéma intermédiaire au schéma de la table wc_out ; cette indépendance facilite les tests et le remplacement de chaque réducteur.

Si vous omettez OutputKeySortColumns , PartitionColumns et OutputGroupingColumns pour un mappeur lors de la construction du pipeline, le framework utilise la OutputKey du mappeur comme valeur par défaut pour ces trois paramètres.

Code Java

package com.aliyun.odps.mapred.open.example;
import java.io.IOException;
import java.util.Iterator;
import com.aliyun.odps.Column;
import com.aliyun.odps.OdpsException;
import com.aliyun.odps.OdpsType;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.data.TableInfo;
import com.aliyun.odps.mapred.Job;
import com.aliyun.odps.mapred.MapperBase;
import com.aliyun.odps.mapred.ReducerBase;
import com.aliyun.odps.pipeline.Pipeline;
public class WordCountPipelineTest {
    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.setBigint(0, 1L);
        }
        @Override
            public void map(long recordNum, Record record, TaskContext context)
            throws IOException {
            for (int i = 0; i < record.getColumnCount(); i++) {
                String[] words = record.get(i).toString().split("\\s+");
                for (String w : words) {
                    word.setString(0, w);
                    context.write(word, one);
                }
            }
        }
    }
    public static class SumReducer extends ReducerBase {
        private Record value;
        @Override
            public void setup(TaskContext context) throws IOException {
            value = context.createOutputValueRecord();
        }
        @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);
            }
            value.set(0, count);
            context.write(key, value);
        }
    }
    public static class IdentityReducer extends ReducerBase {
        private Record result;
        @Override
            public void setup(TaskContext context) throws IOException {
            result = context.createOutputRecord();
        }
        @Override
            public void reduce(Record key, Iterator<Record> values, TaskContext context)
            throws IOException {
            while (values.hasNext()) {
                result.set(0, key.get(0));
                result.set(1, values.next().get(0));
                context.write(result);
            }
        }
    }
    public static void main(String[] args) throws OdpsException {
        if (args.length != 2) {
            System.err.println("Usage: WordCountPipeline <in_table> <out_table>");
            System.exit(2);
        }
        Job job = new Job();
        /** During pipeline construction, if you do not specify OutputKeySortColumns, PartitionColumns, and OutputGroupingColumns for a mapper, the framework uses OutputKey of the mapper as the default values of these parameters.
         */
        Pipeline pipeline = Pipeline.builder()
            .addMapper(TokenizerMapper.class)
            .setOutputKeySchema(
            new Column[] { new Column("word", OdpsType.STRING) })
            .setOutputValueSchema(
            new Column[] { new Column("count", OdpsType.BIGINT) })
            .setOutputKeySortColumns(new String[] { "word" })
            .setPartitionColumns(new String[] { "word" })
            .setOutputGroupingColumns(new String[] { "word" })
            .addReducer(SumReducer.class)
            .setOutputKeySchema(
            new Column[] { new Column("word", OdpsType.STRING) })
            .setOutputValueSchema(
            new Column[] { new Column("count", OdpsType.BIGINT)})
            .addReducer(IdentityReducer.class).createPipeline();
        /** Add the pipeline to jobconf. If you want to configure a combiner, use jobconf. */
        job.setPipeline(pipeline);
        /** Configure the input and output tables. */
        job.addInput(TableInfo.builder().tableName(args[0]).build());
        job.addOutput(TableInfo.builder().tableName(args[1]).build());
        /** Submit the job and wait for it to complete. */
        job.submit();
        job.waitForCompletion();
        System.exit(job.isSuccessful() == true ? 0 : 1);
    }
}

Parcours du flux de données

Avec l'entrée hello,odps, la commande Tunnel charge les données sous forme d'un seul enregistrement dans wc_in : hello est stocké dans la colonne key et odps dans la colonne value.

TokenizerMapper parcourt toutes les colonnes de chaque enregistrement à l'aide de record.getColumnCount(). Pour chaque valeur de colonne, il appelle split("\\s+") et émet une paire (word, 1) par jeton. Comme hello et odps ne contiennent aucun espace, chacun est traité comme un jeton unique. Le mappeur émet :

(hello, 1)
(odps, 1)

SumReducer reçoit toutes les valeurs pour chaque clé de mot et accumule le compte. Chaque mot n'apparaissant qu'une fois, il émet :

(hello, 1)
(odps, 1)

IdentityReducer lit les paires (word, count) issues de SumReducer et les écrit dans la table de sortie wc_out en mappant les positions des colonnes pour correspondre au schéma (key STRING, cnt BIGINT).

Sortie finale dans wc_out :

| hello      | 1          |
| odps       | 1          |