Cet exemple montre comment enchaîner plusieurs jobs MapReduce de manière séquentielle dans MaxCompute. Un job MapReduce ne peut pas en appeler un autre au moment 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 jobs 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 jobs de manière séquentielle :
Job d'initialisation (
InitMapper) : écrit une valeur initiale du compteur (2) dans la table de sortie.Jobs de décrémentation (
DecreaseMapper) : s'exécutent en boucle. À chaque itération, le job 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 se termine lorsque le compteur atteint0.
Le transfert de données entre les jobs repose sur deux mécanismes :
Table de ressources (
multijobs_res_table) : transmet la valeur numérique entre les jobs.Compteur (
multijobs.value) : indique à la méthodemainsi l'itération doit se poursuivre.
La méthode main est le seul orchestrateur : elle soumet les jobs 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 Prise en main
Généré le package
mapreduce-examples.jaret placé celui-ci dans le dossierbin\data\resourcessous le chemin d'installation du client MaxCompute
Préparer les tables et les 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 le job :
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 du job précédent comme une table de ressources. Le fichier JAR est enregistré pour permettre au client MaxCompute de localiser la classe au moment de l'exécution.
Exécuter 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 le job : le package JAR et la table de ressources contenant la sortie du job précédent |
|
|
|
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 à |
Résultat attendu
Une fois le job terminé, 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.");
}
}
}
Parcours du code
InitMapper
S'exécute une fois dans setup(). Lit la valeur initiale du compteur (multijobs.value, par défaut 2) depuis JobConf et l'écrit dans la table de sortie. Il s'agit d'un job de type map uniquement : setNumReduceTasks(0) empêche le framework de lancer des tâches de réduction.
DecreaseMapper
S'exécute dans cleanup() à la fin de chaque itération. Le processus est le suivant :
Lit la valeur attendue (
multijobs.expect.value) depuisJobConf.Lit l'enregistrement unique depuis
multijobs_res_tableen utilisantcontext.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 met à jour le compteur (
multijobs.value) avec la nouvelle valeur.
Méthode main
Contrôle la séquence des jobs :
Soumet le job d'initialisation et attend sa fin.
Entre dans une boucle qui soumet les jobs de décrémentation un par un, en attendant la fin de chacun avant de passer au suivant.
Après chaque job de décrémentation, vérifie le compteur
multijobs.value. Si celui-ci est égal à0, quitte la boucle.Vérifie que
iterCountest égal à0après la boucle ; lève une exceptionIOExceptiondans le cas contraire, signalant un échec du job.