Estes exemplos demonstram como utilizar InputUtils.addTable() com uma especificação de partição para ler dados de partições específicas em um job MapReduce do MaxCompute.
Ambos os exemplos apresentam apenas a função main. O código não está completo e não pode ser compilado ou executado diretamente — utilize-o como referência ao criar sua própria implementação.
Exemplo 1: Leitura de uma única partição
Adote este padrão quando o valor da partição for conhecido no momento do envio do 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);
}
Exemplo 2: Leitura dinâmica de múltiplas partições
Utilize esta abordagem quando for necessário filtrar partições durante a execução. Este exemplo combina o SDK do MaxCompute e o SDK do MapReduce: o SDK do MaxCompute lista todas as partições da tabela, enquanto uma função applicable personalizada define quais partições incluir como entrada.
A função applicable contém a lógica personalizada que você implementa para filtrar partições conforme seus requisitos.
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);
}
Funcionamento do loop de partições:
odps.tables().get(tblname)utiliza o SDK do MaxCompute para recuperar os metadados da tabela, incluindo todas as partições.O loop
forpercorre cada partição e invocaapplicable(p)— uma função personalizada implementada por você — para decidir se ela deve ser incluída.Para cada partição incluída, um objeto
LinkedHashMap<String, String>é construído a partir dos pares chave-valor da partição e depois passado paraInputUtils.addTable()como especificação de partição.A tabela de saída é adicionada sem especificação de partição por meio de
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job).