Cet exemple montre comment enregistrer un package JAR et un fichier texte en tant que ressources MaxCompute, puis les utiliser dans une tâche MapReduce pour charger des données dans une table.
Fonctionnement
Lors de la phase de configuration, le mapper lit le fichier texte à l'aide de context.readResourceFileAsStream(), analyse chaque ligne et écrit les enregistrements dans la table de sortie. La tâche ne comporte aucune phase de réduction.
Les options -resources et -classpath ont des fonctions distinctes :
|
Option |
Objectif |
|
|
Enregistre les fichiers afin que MaxCompute les distribue aux nœuds de calcul. Utilisez cette option pour tout fichier que la tâche doit lire lors de l'exécution. |
|
|
Ajoute le fichier JAR au classpath Java afin que la JVM puisse charger la classe |
Prérequis
Avant de commencer, assurez-vous d'avoir :
Effectué la configuration de l'environnement décrite dans la section Démarrage rapide
Le fichier
mapreduce-examples.jarprésent dans le répertoirebin\data\resourcesde votre installation locale MaxCompute
Préparer les tables et les ressources
Exécutez les commandes suivantes sur le client MaxCompute.
-
Créez la table de sortie :
CREATE TABLE mr_upload_src(key BIGINT, value STRING); -
Ajoutez le fichier texte et le package JAR en tant que ressources :
add file data\resources\import.txt -f; add jar data\resources\mapreduce-examples.jar -f;Lors du premier ajout du package JAR, ignorez l'option
-f.Le fichier
import.txtcontient les données suivantes :1000,odps
Exécuter la tâche
Exécutez la commande suivante sur le client MaxCompute pour démarrer la tâche de chargement :
jar -resources mapreduce-examples.jar,import.txt -classpath data\resources\mapreduce-examples.jar
com.aliyun.odps.mapred.open.example.Upload import.txt mr_upload_src;
Référence des paramètres :
|
Paramètre |
Description |
|
|
Enregistre le fichier JAR et le fichier texte afin que MaxCompute les distribue aux nœuds de calcul |
|
|
Ajoute le fichier JAR au classpath Java afin que la JVM puisse charger la classe |
|
|
Premier argument transmis à |
|
|
Deuxième argument transmis à |
Vérifier le résultat
Une fois la tâche terminée, interrogez la table mr_upload_src. La table doit contenir la ligne suivante, qui correspond directement à la ligne 1000,odps du fichier import.txt (champ 0 → key, champ 1 → value) :
+------------+------------+
| key | value |
+------------+------------+
| 1000 | odps |
+------------+------------+
Exemple de code
Pour la configuration des dépendances du modèle objet de projet (POM), consultez la section Démarrage rapide.
package com.aliyun.odps.mapred.open.example;
import java.io.BufferedInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
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.TaskContext;
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;
/**
* Upload
* Import data from a text file into a table.
*/
public class Upload {
public static class UploadMapper extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
Record record = context.createOutputRecord();
StringBuilder importdata = new StringBuilder();
BufferedInputStream bufferedInput = null;
try {
byte[] buffer = new byte[1024];
int bytesRead = 0;
// Get the resource file name from JobConf
String filename = context.getJobConf().get("import.filename");
// Read the resource file distributed to this worker node
bufferedInput = context.readResourceFileAsStream(filename);
while ((bytesRead = bufferedInput.read(buffer)) != -1) {
String chunk = new String(buffer, 0, bytesRead);
importdata.append(chunk);
}
// Parse each line: split by comma, write key (BIGINT) and value (STRING)
String lines[] = importdata.toString().split("\n");
for (int i = 0; i < lines.length; i++) {
String[] ss = lines[i].split(",");
record.set(0, Long.parseLong(ss[0].trim()));
record.set(1, ss[1].trim());
context.write(record);
}
} catch (FileNotFoundException ex) {
throw new IOException(ex);
} catch (IOException ex) {
throw new IOException(ex);
} finally {
}
}
@Override
public void map(long recordNum, Record record, TaskContext context)
throws IOException {
}
}
public static void main(String[] args) throws Exception {
if (args.length != 2) {
System.err.println("Usage: Upload <import_txt> <out_table>");
System.exit(2);
}
JobConf job = new JobConf();
job.setMapperClass(UploadMapper.class);
// Pass the resource file name to the mapper via JobConf
job.set("import.filename", args[0]);
// Set reducers to 0: this is a map-only job
job.setNumReduceTasks(0);
job.setMapOutputKeySchema(SchemaUtils.fromString("key:bigint"));
job.setMapOutputValueSchema(SchemaUtils.fromString("value:string"));
InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), job);
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
JobClient.runJob(job);
}
}
Configurer JobConf
L'exemple ci-dessus utilise l'interface JobConf du SDK pour définir les propriétés de la tâche. Vous pouvez également utiliser le paramètre -conf de la commande jar pour spécifier un fichier de configuration JobConf :
jar -conf <your-jobconf-file> ...