Dans une tâche de type map-only, le mapper écrit les enregistrements de sortie directement dans une table MaxCompute ; aucun reducer n'est exécuté. Contrairement au MapReduce standard, vous devez uniquement spécifier les tables de sortie, sans définir de métadonnées clé-valeur pour la sortie du mapper.
Cet exemple illustre trois aspects :
Configurer une tâche map-only en définissant le nombre de reducers à 0
Transmettre des paramètres via JobConf et les lire au sein du mapper
Comprendre 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;La commande charge les lignes suivantes dans la table
wc_in:hello,odps hello,odps
Exécution de la tâche
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 de la tâche |
|
|
Spécifie le chemin d'accès au package JAR |
|
|
Table d'entrée |
|
|
Table de sortie |
|
|
Définit |
Résultat attendu
Une fois la tâche terminée, interrogez la table wc_out :
+------------+------------+
| key | cnt |
+------------+------------+
| hello | 1 |
| hello | 1 |
+------------+------------+
La table contient deux lignes au lieu d'une, car une tâche 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 de la tâche
L'instruction job.setNumReduceTasks(0) est requise pour les tâches map-only. Sans elle, le framework attend un reducer et la tâche échoue. Les méthodes InputUtils.addTable et OutputUtils.addTable associent les tables d'entrée et de sortie à la tâche 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() — s'exécute 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() — s'exécute 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 par ligne d'entrée.
cleanup() — s'exécute 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.