Cet exemple montre comment enchaîner plusieurs tâches MapReduce de manière séquentielle dans MaxCompute. Une tâche MapReduce ne peut pas en invoquer une autre lors de l'exécution. Par conséquent, si votre logique de traitement nécessite une exécution itérative — où chaque itération dépend de la sortie de la précédente —, vous devez orchestrer plusieurs tâches depuis une méthode main côté client. Cet exemple utilise un compteur pour contrôler la boucle d'itération.
Fonctionnement
Le programme exécute deux types de tâches de manière séquentielle :
Tâche d'initialisation (
InitMapper) : Écrit une valeur initiale du compteur (2) dans la table de sortie.Tâches de décrémentation (
DecreaseMapper) : S'exécutent en boucle. À chaque itération, le système lit la sortie précédente depuis une table de ressources, décrémente la valeur de 1 et met à jour un compteur. La boucle s'arrête lorsque le compteur atteint0.
Le transfert de données entre les tâches repose sur deux mécanismes :
Table de ressources (
multijobs_res_table) : Transmet la valeur numérique d'une tâche à l'autre.Compteur (
multijobs.value) : Indique à la méthodemains'il faut poursuivre l'itération.
La méthode main est le seul orchestrateur : elle soumet les tâches dans l'ordre et vérifie leur achèvement entre chaque étape.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Effectué la configuration de l'environnement décrite dans la section Prise en main
Généré le package
mapreduce-examples.jaret placé celui-ci dans le répertoirebin\data\resourcesdu chemin d'installation du client MaxCompute
Préparation des tables et des ressources
-
Créez les tables de test :
CREATE TABLE mr_empty (key STRING, value STRING); CREATE TABLE mr_multijobs_out (value BIGINT); -
Enregistrez les ressources utilisées par la tâche :
add table mr_multijobs_out as multijobs_res_table -f; -- Omit -f when adding the JAR for the first time. add jar data\resources\mapreduce-examples.jar -f;La table
mr_multijobs_outest enregistrée en tant quemultijobs_res_tableafin queDecreaseMapperpuisse lire la sortie de la tâche précédente comme une table de ressources. Le fichier JAR est enregistré pour permettre au client MaxCompute de localiser la classe lors de l'exécution.
Exécution de MultiJobs
Exécutez la commande suivante sur le client MaxCompute :
jar -resources mapreduce-examples.jar,multijobs_res_table -classpath data\resources\mapreduce-examples.jar \
com.aliyun.odps.mapred.open.example.MultiJobs mr_multijobs_out;
Paramètres de la commande :
|
Paramètre |
Valeur |
Description |
|
|
|
Ressources disponibles pour la tâche : le package JAR et la table de ressources contenant la sortie de la tâche précédente |
|
|
|
Chemin local du fichier JAR utilisé par le client pour localiser la classe principale |
|
Classe principale |
|
Point d'entrée du programme |
|
Argument |
|
Nom de la table de sortie transmis à la méthode |
Résultat attendu
Une fois la tâche terminée, la table mr_multijobs_out contient un seul enregistrement :
+------------+
| value |
+------------+
| 0 |
+------------+
Exemple de code
Pour la configuration des dépendances du modèle objet de projet (POM), consultez la section Précautions du guide de prise en main.
package com.aliyun.odps.mapred.open.example;
import java.io.IOException;
import java.util.Iterator;
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.RunningJob;
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;
/**
* MultiJobs
*
* Running multiple job
*
**/
public class MultiJobs {
public static class InitMapper extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
Record record = context.createOutputRecord();
long v = context.getJobConf().getLong("multijobs.value", 2);
record.set(0, v);
context.write(record);
}
}
public static class DecreaseMapper extends MapperBase {
@Override
public void cleanup(TaskContext context) throws IOException {
/** Obtain the variable values that are defined in the main function from JobConf. */
long expect = context.getJobConf().getLong("multijobs.expect.value", -1);
long v = -1;
int count = 0;
/** Read the data from the output table of the previous job. */
Iterator<Record> iter = context.readResourceTable("multijobs_res_table");
while (iter.hasNext()) {
Record r = iter.next();
v = (Long) r.get(0);
if (expect != v) {
throw new IOException("expect: " + expect + ", but: " + v);
}
count++;
}
if (count != 1) {
throw new IOException("res_table should have 1 record, but: " + count);
}
Record record = context.createOutputRecord();
v--;
record.set(0, v);
context.write(record);
/** Set the counter. The counter value can be obtained in the main function after the job is completed. */
context.getCounter("multijobs", "value").setValue(v);
}
}
public static void main(String[] args) throws Exception {
if (args.length != 1) {
System.err.println("Usage: TestMultiJobs <table>");
System.exit(1);
}
String tbl = args[0];
long iterCount = 2;
System.err.println("Start to run init job.");
JobConf initJob = new JobConf();
initJob.setLong("multijobs.value", iterCount);
initJob.setMapperClass(InitMapper.class);
InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), initJob);
OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), initJob);
initJob.setMapOutputKeySchema(SchemaUtils.fromString("key:string"));
initJob.setMapOutputValueSchema(SchemaUtils.fromString("value:string"));
/** Explicitly set the number of reducers to 0 for map-only jobs. */
initJob.setNumReduceTasks(0);
JobClient.runJob(initJob);
while (true) {
System.err.println("Start to run iter job, count: " + iterCount);
JobConf decJob = new JobConf();
decJob.setLong("multijobs.expect.value", iterCount);
decJob.setMapperClass(DecreaseMapper.class);
InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), decJob);
OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), decJob);
/** Explicitly set the number of reducers to 0 for map-only jobs. */
decJob.setNumReduceTasks(0);
RunningJob rJob = JobClient.runJob(decJob);
iterCount--;
/** If the specified number of iterations is reached, exit the loop. */
if (rJob.getCounters().findCounter("multijobs", "value").getValue() == 0) {
break;
}
}
if (iterCount != 0) {
throw new IOException("Job failed.");
}
}
}
Analyse du code
InitMapper
S'exécute une seule fois dans la méthode setup(). Elle lit la valeur initiale du compteur (multijobs.value, par défaut 2) depuis l'objet JobConf et l'écrit dans la table de sortie. Il s'agit d'une tâche de type map uniquement ; l'appel à setNumReduceTasks(0) empêche le framework de lancer des réducteurs.
DecreaseMapper
S'exécute dans la méthode cleanup() à la fin de chaque itération. Cette classe effectue les opérations suivantes :
Lit la valeur attendue (
multijobs.expect.value) depuis l'objetJobConf.Lit l'unique enregistrement de la table
multijobs_res_tableà l'aide de la méthodecontext.readResourceTable(). Lève une exceptionIOExceptionsi le nombre d'enregistrements n'est pas exactement égal à 1 ou si la valeur ne correspond pas à la valeur attendue.Décrémente la valeur de 1, l'écrit dans la table de sortie et définit le compteur (
multijobs.value) avec la nouvelle valeur.
Méthode main
Contrôle la séquence des tâches :
Soumet la tâche d'initialisation et attend son achèvement.
Entre dans une boucle qui soumet les tâches de décrémentation une par une, en attendant la fin de chacune avant de passer à la suivante.
Après chaque tâche de décrémentation, vérifie la valeur du compteur
multijobs.value. Si celle-ci est égale à0, la boucle se termine.Vérifie que la variable
iterCountest bien égale à0après la boucle ; sinon, lève une exceptionIOExceptionpour signaler l'échec de la tâche.