SecondarySort trie la sortie MapReduce à l'aide d'une clé composite : il regroupe les enregistrements par une clé primaire tout en triant les valeurs selon une clé secondaire au sein de chaque groupe. Cet exemple montre comment implémenter SecondarySort dans MaxCompute MapReduce en utilisant trois paramètres de job qui fonctionnent conjointement : les colonnes de tri, les colonnes de partitionnement et les colonnes de regroupement.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Effectué la configuration de l'environnement décrite dans Prise en main
Fonctionnement
Le mapper lit chaque enregistrement d'entrée contenant deux entiers et émet une paire clé-valeur composite ((key, value), value). La clé composite contient les deux entiers afin que le framework puisse trier sur les deux champs.
Trois paramètres de job contrôlent la suite des opérations :
|
Paramètre |
Colonnes |
Objectif |
|
Colonnes de tri |
|
Trie toute la sortie du map d'abord par clé primaire, puis par clé secondaire au sein de chaque groupe de clés primaires |
|
Colonnes de partitionnement |
|
Achemine les enregistrements ayant la même clé primaire vers le même reducer |
|
Colonnes de regroupement |
|
Appelle la méthode |
Le reducer reçoit un groupe par clé primaire et itère sur ses valeurs, déjà triées. Il écrit chaque paire (key, value) dans la table de sortie.
Préparation des tables et des ressources
1. Créez les tables d'entrée et de sortie.
CREATE TABLE ss_in(key BIGINT, value BIGINT);
CREATE TABLE ss_out(key BIGINT, value BIGINT);
2. Ajoutez le package JAR en tant que ressource MaxCompute.
add jar data\resources\mapreduce-examples.jar -f;
Omettez l'indicateur -f lors de l'ajout du package JAR pour la première fois.
3. Téléchargez les données de test vers ss_in à l'aide de Tunnel.
tunnel upload data.txt ss_in;
Le fichier data.txt se trouve dans le répertoire bin du client MaxCompute et contient :
1,2
2,1
1,1
2,2
Exécution de SecondarySort
Exécutez la commande suivante sur le client MaxCompute :
jar -resources mapreduce-examples.jar -classpath data\resources\mapreduce-examples.jar com.aliyun.odps.mapred.open.example.SecondarySort ss_in ss_out;
Sortie attendue
Une fois le job terminé, interrogez ss_out. Les enregistrements sont triés par key croissant, puis par value croissant au sein de chaque groupe de clés :
+------------+------------+
| key | value |
+------------+------------+
| 1 | 1 |
| 1 | 2 |
| 2 | 1 |
| 2 | 2 |
+------------+------------+
Parcours du flux de données
La section suivante illustre la transformation des données à chaque étape.
Enregistrements d'entrée (ss_in) :
(1, 2)
(2, 1)
(1, 1)
(2, 2)
Sortie du map — clé composite (i1, i2) associée à la valeur i2 :
key=(1,1), value=1
key=(1,2), value=2
key=(2,1), value=1
key=(2,2), value=2
Le framework trie toute la sortie du map par (i1, i2) et achemine les enregistrements par i1 vers le reducer approprié.
Entrée du reduce — deux groupes, un pour chaque valeur unique de i1 :
Group i1=1: values [1, 2]
Group i1=2: values [1, 2]
Sortie (ss_out) :
(1, 1), (1, 2), (2, 1), (2, 2)
Exemple de code
Pour la configuration des dépendances du Project Object Model (POM), consultez les Précautions.
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.mapred.JobClient;
import com.aliyun.odps.mapred.MapperBase;
import com.aliyun.odps.mapred.ReducerBase;
import com.aliyun.odps.mapred.TaskContext;
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;
/**
*
* This is an example ODPS Map/Reduce application. It reads the input table that
* must contain two integers per record. The output is sorted by the first and
* second number and grouped on the first number.
*
**/
public class SecondarySort {
/**
* Read two integers from each line and generate a key, value pair as ((left,
* right), right).
**/
public static class MapClass extends MapperBase {
private Record key;
private Record value;
@Override
public void setup(TaskContext context) throws IOException {
key = context.createMapOutputKeyRecord();
value = context.createMapOutputValueRecord();
}
@Override
public void map(long recordNum, Record record, TaskContext context)
throws IOException {
long left = 0;
long right = 0;
if (record.getColumnCount() > 0) {
left = (Long) record.get(0);
if (record.getColumnCount() > 1) {
right = (Long) record.get(1);
}
key.set(new Object[] { (Long) left, (Long) right });
value.set(new Object[] { (Long) right });
context.write(key, value);
}
}
}
/**
* A reducer class that just emits the sum of the input values.
**/
public static class ReduceClass 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 {
result.set(0, key.get(0));
while (values.hasNext()) {
Record value = values.next();
result.set(1, value.get(0));
context.write(result);
}
}
}
public static void main(String[] args) throws Exception {
if (args.length != 2) {
System.err.println("Usage: secondarysrot <in> <out>");
System.exit(2);
}
JobConf job = new JobConf();
job.setMapperClass(MapClass.class);
job.setReducerClass(ReduceClass.class);
/** Set the composite key columns that determine sort order. */
// Sort by i1 first, then by i2 within each i1 group
job.setOutputKeySortColumns(new String[] { "i1", "i2" });
// Route records with the same i1 value to the same reducer
job.setPartitionColumns(new String[] { "i1" });
// Call reduce() once per unique i1 value
job.setOutputGroupingColumns(new String[] { "i1" });
// The map output key is a pair (i1, i2); the value carries i2
job.setMapOutputKeySchema(SchemaUtils.fromString("i1:bigint,i2:bigint"));
job.setMapOutputValueSchema(SchemaUtils.fromString("i2x:bigint"));
InputUtils.addTable(TableInfo.builder().tableName(args[0]).build(), job);
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
JobClient.runJob(job);
System.exit(0);
}
}