SecondarySort trie la sortie MapReduce à l'aide d'une clé composite : il regroupe les enregistrements selon une clé primaire tout en triant les valeurs par une clé secondaire au sein de chaque groupe. Cet exemple illustre l'implémentation de SecondarySort dans MaxCompute MapReduce grâce à trois paramètres de job qui agissent conjointement : les colonnes de tri, les colonnes de partitionnement et les colonnes de regroupement.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Configuré l'environnement comme décrit 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, ce qui permet au framework de trier selon ces deux critères.
Trois paramètres de job contrôlent la suite du traitement :
|
Paramètre |
Colonnes |
Objectif |
|
Sort columns |
|
Trie toute la sortie du map d'abord par la clé primaire, puis par la clé secondaire au sein de chaque groupe de clés primaires |
|
Partition columns |
|
Achemine les enregistrements partageant la même clé primaire vers le même reducer |
|
Grouping columns |
|
Appelle la méthode |
Le reducer reçoit un groupe par clé primaire et parcourt 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. Chargez les données de test dans 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 la table 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 montre comment les données d'entrée sont transformé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 selon (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 la section 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);
}
}