Ces exemples montrent comment utiliser la méthode InputUtils.addTable() avec une spécification de partition pour lire des données issues de partitions spécifiques dans un job MapReduce MaxCompute.
Les deux exemples présentent uniquement la fonction main. Le code est incomplet et ne peut ni être compilé ni exécuté directement. Utilisez-le comme référence lors de la mise en œuvre de votre propre solution.
Exemple 1 : Lecture depuis une seule partition
Utilisez ce modèle lorsque la valeur de la partition est connue au moment de la soumission du job.
public static void main(String[] args) throws Exception {
JobConf job = new JobConf();
...
LinkedHashMap<String, String> input = new LinkedHashMap<String, String>();
input.put("pt", "123456");
InputUtils.addTable(TableInfo.builder().tableName("input_table").partSpec(input).build(), job);
LinkedHashMap<String, String> output = new LinkedHashMap<String, String>();
output.put("ds", "654321");
OutputUtils.addTable(TableInfo.builder().tableName("output_table").partSpec(output).build(), job);
JobClient.runJob(job);
}
Exemple 2 : Lecture dynamique depuis plusieurs partitions
Privilégiez cette approche si vous devez filtrer les partitions à l'exécution. Cet exemple combine le SDK MaxCompute et le SDK MapReduce : le SDK MaxCompute répertorie toutes les partitions de la table, tandis qu'une fonction personnalisée applicable détermine les partitions à inclure en entrée.
La fonction applicable correspond à une logique personnalisée que vous implémentez pour filtrer les partitions selon vos besoins.
package com.aliyun.odps.mapred.open.example;
...
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);
job.setMapOutputKeySchema(SchemaUtils.fromString("word:string"));
job.setMapOutputValueSchema(SchemaUtils.fromString("count:bigint"));
// Using an Alibaba Cloud account's AccessKey pair grants access to all API operations,
// which is a high-risk approach. Use a RAM user instead for routine operations.
// To create a RAM user, go to the Resource Access Management (RAM) console.
// Store credentials in environment variables rather than hardcoding them in your code.
Account account = new AliyunAccount(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
Odps odps = new Odps(account);
odps.setEndpoint("odps_endpoint_url");
odps.setDefaultProject("my_project");
Table table = odps.tables().get(tblname);
TableInfoBuilder builder = TableInfo.builder().tableName(tblname);
for (Partition p : table.getPartitions()) {
if (applicable(p)) {
LinkedHashMap<String, String> partSpec = new LinkedHashMap<String, String>();
for (String key : p.getPartitionSpec().keys()) {
partSpec.put(key, p.getPartitionSpec().get(key));
}
InputUtils.addTable(builder.partSpec(partSpec).build(), job);
}
}
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
JobClient.runJob(job);
}
Fonctionnement de la boucle sur les partitions :
La méthode
odps.tables().get(tblname)utilise le SDK MaxCompute pour récupérer les métadonnées de la table, y compris toutes les partitions.La boucle
forparcourt chaque partition et appelleapplicable(p)— une fonction personnalisée que vous implémentez — pour décider de son inclusion.Pour chaque partition retenue, un objet
LinkedHashMap<String, String>est construit à partir des paires clé-valeur de la partition, puis transmis àInputUtils.addTable()en tant que spécification de partition.La table de sortie est ajoutée sans spécification de partition via l'appel
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job).