Cet exemple vous guide pas à pas dans l'exécution d'une tâche WordCount avec MapReduce sur MaxCompute, de la création des tables d'entrée et du chargement des données jusqu'à l'exécution de la tâche et la vérification des résultats.
Fonctionnement
L'exemple WordCount suit le flux de données standard de MapReduce :
(input) <k1, v1> -> map -> <k2, v2> -> combine -> <k2, v2> -> reduce -> <k3, v3> (output)
Trois classes implémentent ce pipeline :
TokenizerMapper : lit chaque enregistrement de la table d'entrée et émet chaque valeur de champ sous forme de clé associée au nombre
1SumCombiner : agrège les comptes localement sur le nœud mapper avant d'envoyer les données au réducteur, ce qui réduit le trafic réseau
SumReducer : additionne tous les comptes pour chaque clé unique et écrit le résultat final dans la table de sortie
Prérequis
Avant de commencer, assurez-vous d'avoir :
Effectué la configuration de l'environnement décrite dans la section Prise en main
Préparation de l'environnement
1. Préparez le package JAR.
Le package JAR de cet exemple est mapreduce-examples.jar, situé dans le répertoire bin\data\resources du chemin d'installation du client MaxCompute.
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 -f lors de l'ajout du package JAR pour la première fois.
2. Chargez les données d'entrée.
Utilisez Tunnel pour charger le fichier data.txt depuis le répertoire bin du client MaxCompute vers la table wc_in :
tunnel upload data.txt wc_in;
Les données suivantes sont chargées dans wc_in :
hello,odps
Exécution de WordCount
Exécutez la tâche WordCount sur le client MaxCompute :
jar -resources mapreduce-examples.jar -classpath data\resources\mapreduce-examples.jar com.aliyun.odps.mapred.open.example.WordCount wc_in wc_out
Vérification des résultats
Une fois la tâche terminée, interrogez la table wc_out pour afficher les résultats :
SELECT * FROM wc_out;
Le résultat doit correspondre à ce qui suit :
+------------+------------+
| key | cnt |
+------------+------------+
| hello | 1 |
| odps | 1 |
+------------+------------+
Exemple de code
Pour la configuration des dépendances du modèle objet de projet (POM), consultez la section Précautions de la rubrique Prise en main.Usage notes
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.ReducerBase;
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;
public class WordCount {
// Mapper: reads each record, emits (word, 1) for every field value
public static class TokenizerMapper extends MapperBase {
private Record word;
private Record one;
@Override
public void setup(TaskContext context) throws IOException {
word = context.createMapOutputKeyRecord();
one = context.createMapOutputValueRecord();
one.set(new Object[] { 1L });
System.out.println("TaskID:" + context.getTaskID().toString());
}
@Override
public void map(long recordNum, Record record, TaskContext context)
throws IOException {
for (int i = 0; i < record.getColumnCount(); i++) {
word.set(new Object[] { record.get(i).toString() });
context.write(word, one);
}
}
}
// Combiner: aggregates map output locally before shuffling to the reducer
public static class SumCombiner extends ReducerBase {
private Record count;
@Override
public void setup(TaskContext context) throws IOException {
count = context.createMapOutputValueRecord();
}
@Override
public void reduce(Record key, Iterator<Record> values, TaskContext context)
throws IOException {
long c = 0;
while (values.hasNext()) {
Record val = values.next();
c += (Long) val.get(0);
}
count.set(0, c);
context.write(key, count);
}
}
// Reducer: sums counts for each unique key and writes the final output
public static class SumReducer extends ReducerBase {
private Record result = null;
@Override
public void setup(TaskContext context) throws IOException {
result = context.createOutputRecord();
}
@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);
}
result.set(0, key.get(0));
result.set(1, count);
context.write(result);
}
}
public static void main(String[] args) throws Exception {
if (args.length != 2) {
System.err.println("Usage: WordCount <in_table> <out_table>");
System.exit(2);
}
JobConf job = new JobConf();
job.setMapperClass(TokenizerMapper.class);
job.setCombinerClass(SumCombiner.class);
job.setReducerClass(SumReducer.class);
// Define the intermediate key-value schema between mapper and reducer
job.setMapOutputKeySchema(SchemaUtils.fromString("word:string"));
job.setMapOutputValueSchema(SchemaUtils.fromString("count:bigint"));
// Set input and output tables
InputUtils.addTable(TableInfo.builder().tableName(args[0]).build(), job);
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
JobClient.runJob(job);
}
}
Parcours détaillé
Pour illustrer le flux de données à travers le pipeline, voici ce que produit chaque étape pour l'entrée hello,odps :
Sortie du mapper — TokenizerMapper émet une paire clé-valeur par champ :
<hello, 1>
<odps, 1>
Sortie du combiner — SumCombiner effectue une agrégation locale (dans ce petit exemple, les comptes restent à 1) :
<hello, 1>
<odps, 1>
Sortie du réducteur — SumReducer écrit les totaux finaux dans wc_out :
<hello, 1>
<odps, 1>