Dans un job de type « map-only », le mapper écrit directement les enregistrements de sortie dans une table MaxCompute : aucun reducer n'est exécuté. Contrairement au modèle MapReduce standard, il suffit de spécifier les tables de sortie, sans définir les métadonnées clé-valeur pour la sortie du mapper.
Cet exemple illustre trois aspects :
La configuration d'un job map-only en définissant le nombre de reducers à 0
Le passage de paramètres via JobConf et leur lecture au sein du mapper
Le fonctionnement des méthodes de cycle de vie
setup,mapetcleanupavec une exécution conditionnelle
Prérequis
Avant de commencer, configurez l'environnement comme décrit dans la rubrique Prise en main.
Préparation des tables de test et des ressources
-
Créez les tables d'entrée et de sortie.
CREATE TABLE wc_in (key STRING, value STRING); CREATE TABLE wc_out (key STRING, cnt BIGINT); -
Ajoutez le package JAR en tant que ressource.
add jar data\resources\mapreduce-examples.jar -f;Omettez l'option
-flors du premier ajout du package JAR. Le chemindata\resources\mapreduce-examples.jarest relatif au répertoirebinde votre installation locale du client MaxCompute. -
Importez les données de test dans la table
wc_inà l'aide de Tunnel. Exécutez la commande suivante depuis le répertoirebindu client MaxCompute, où se trouve le fichierdata.txt.tunnel upload data.txt wc_in;Cette commande charge les lignes suivantes dans la table
wc_in:hello,odps hello,odps
Exécution du job
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.MapOnly wc_in wc_out map
Les arguments de la commande correspondent aux éléments suivants :
|
Argument |
Description |
|
|
Déclare le package JAR comme dépendance du job |
|
|
Spécifie le chemin d'accès au package JAR |
|
|
Table d'entrée |
|
|
Table de sortie |
|
|
Définit |
Résultat attendu
Une fois le job terminé, interrogez la table wc_out :
+------------+------------+
| key | cnt |
+------------+------------+
| hello | 1 |
| hello | 1 |
+------------+------------+
La table contient deux lignes au lieu d'une, car un job map-only produit un enregistrement de sortie pour chaque enregistrement d'entrée, sans agrégation. Chaque ligne d'entrée hello,odps correspond à une ligne de sortie hello | 1.
Exemple de code
Pour les dépendances du Project Object Model (POM), consultez la section Précautions.
package com.aliyun.odps.mapred.open.example;
import java.io.IOException;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.mapred.JobClient;
import com.aliyun.odps.mapred.MapperBase;
import com.aliyun.odps.mapred.conf.JobConf;
import com.aliyun.odps.mapred.utils.SchemaUtils;
import com.aliyun.odps.mapred.utils.InputUtils;
import com.aliyun.odps.mapred.utils.OutputUtils;
import com.aliyun.odps.data.TableInfo;
public class MapOnly {
public static class MapperClass extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.setup", false);
/** The main function executes the following logic only if option.mapper.setup is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, "setup");
result.set(1, 1L);
context.write(result);
}
}
@Override
public void map(long key, Record record, TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.map", false);
/** The main function executes the following logic only if option.mapper.map is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, record.get(0));
result.set(1, 1L);
context.write(result);
}
}
@Override
public void cleanup(TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.cleanup", false);
/** The main function executes the following logic only if option.mapper.cleanup is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, "cleanup");
result.set(1, 1L);
context.write(result);
}
}
}
public static void main(String[] args) throws Exception {
if (args.length != 2 && args.length != 3) {
System.err.println("Usage: OnlyMapper <in_table> <out_table> [setup|map|cleanup]");
System.exit(2);
}
JobConf job = new JobConf();
job.setMapperClass(MapperClass.class);
/** For MapOnly jobs, the number of reducers must be explicitly set to 0. */
job.setNumReduceTasks(0);
/** Configure information about input and output tables. */
InputUtils.addTable(TableInfo.builder().tableName(args[0]).build(), job);
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
if (args.length == 3) {
String options = new String(args[2]);
/** You can specify key-value pairs in the JobConf file, and use getJobConf of the context to query the configurations in a mapper. */
if (options.contains("setup")) {
job.setBoolean("option.mapper.setup", true);
}
if (options.contains("map")) {
job.setBoolean("option.mapper.map", true);
}
if (options.contains("cleanup")) {
job.setBoolean("option.mapper.cleanup", true);
}
}
JobClient.runJob(job);
}
}
Fonctionnement du code
main() — configuration du job
L'instruction job.setNumReduceTasks(0) est requise pour les jobs map-only. Sans elle, le framework s'attend à trouver un reducer et le job échoue. Les méthodes InputUtils.addTable et OutputUtils.addTable associent les tables d'entrée et de sortie au job en utilisant les arguments de la ligne de commande.
Le troisième argument facultatif (setup, map ou cleanup) définit l'indicateur booléen correspondant dans JobConf. Chaque méthode de cycle de vie lit son propre indicateur via context.getJobConf().getBoolean(...) et n'écrit la sortie que si l'indicateur est défini sur true.
setup() — exécutée une fois avant le début du traitement
Écrit un seul enregistrement avec la clé "setup" et le compte 1. S'exécute lorsque option.mapper.setup=true.
map() — exécutée une fois par enregistrement d'entrée
Lit le premier champ de chaque enregistrement d'entrée (record.get(0)), et l'écrit avec le compte 1. S'exécute lorsque option.mapper.map=true. Dans cet exemple, l'argument map active cette méthode, ce qui explique pourquoi la sortie contient deux lignes, soit une pour chaque ligne d'entrée.
cleanup() — exécutée une fois après le traitement de tous les enregistrements
Écrit un seul enregistrement avec la clé "cleanup" et le compte 1. S'exécute lorsque option.mapper.cleanup=true.