Em um job map-only, o mapper grava os registros de saída diretamente em uma tabela do MaxCompute, sem executar nenhum reducer. Diferentemente do MapReduce padrão, basta especifique as tabelas de saída; não é necessário defina metadados de chave-valor para a saída do mapper.
Este exemplo demonstra três aspectos:
Configure de um job map-only ao defina a contagem de reducers como 0
Passagem de parâmetros via JobConf e leitura dentro do mapper
Funcionamento dos métodos de ciclo de vida
setup,mapecleanupcom execução condicional
Pré-requisitos
Antes de começar, conclua a configure do ambiente descrita em Primeiros passos.
Prepare tabelas e recursos de teste
-
Crie as tabelas de entrada e de saída.
CREATE TABLE wc_in (key STRING, value STRING); CREATE TABLE wc_out (key STRING, cnt BIGINT); -
Adicione o pacote JAR como recurso.
add jar data\resources\mapreduce-examples.jar -f;Omita
-fna primeira vez que adicionar o pacote JAR. O caminhodata\resources\mapreduce-examples.jaré relativo ao diretóriobinda instalação local do cliente maxcompute. -
Importe dados de teste para
wc_inusando o Tunnel. Execute o comando abaixo no diretóriobindo cliente maxcompute, onde está o arquivodata.txt.tunnel upload data.txt wc_in;O comando carrega as seguintes linhas em
wc_in:hello,odps hello,odps
Execute o job
Execute o seguinte comando no cliente maxcompute:
jar -resources mapreduce-examples.jar -classpath data\resources\mapreduce-examples.jar
com.aliyun.odps.mapred.open.example.MapOnly wc_in wc_out map
Os argumentos do comando têm o seguinte mapeamento:
|
Argumento |
Descrição |
|
|
Declara o pacote JAR como dependência do job |
|
|
Especifique o caminho para o pacote JAR |
|
|
Tabela de entrada |
|
|
Tabela de saída |
|
|
Defina |
Resultado esperado
Após a conclusão do job, consulte wc_out:
+------------+------------+
| key | cnt |
+------------+------------+
| hello | 1 |
| hello | 1 |
+------------+------------+
A tabela contém duas linhas em vez de uma porque um job map-only produz um registro de saída para cada registro de entrada, sem agregação. Cada linha de entrada hello,odps corresponde a uma linha de saída hello | 1.
Código de exemplo
Para dependências do Project Object Model (POM), consulte Precauções.
package com.aliyun.odps.mapred.open.example;
import java.io.IOException;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.mapred.JobClient;
import com.aliyun.odps.mapred.MapperBase;
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;
public class MapOnly {
public static class MapperClass extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.setup", false);
/** The main function executes the following logic only if option.mapper.setup is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, "setup");
result.set(1, 1L);
context.write(result);
}
}
@Override
public void map(long key, Record record, TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.map", false);
/** The main function executes the following logic only if option.mapper.map is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, record.get(0));
result.set(1, 1L);
context.write(result);
}
}
@Override
public void cleanup(TaskContext context) throws IOException {
boolean is = context.getJobConf().getBoolean("option.mapper.cleanup", false);
/** The main function executes the following logic only if option.mapper.cleanup is set to true in the JobConf file: */
if (is) {
Record result = context.createOutputRecord();
result.set(0, "cleanup");
result.set(1, 1L);
context.write(result);
}
}
}
public static void main(String[] args) throws Exception {
if (args.length != 2 && args.length != 3) {
System.err.println("Usage: OnlyMapper <in_table> <out_table> [setup|map|cleanup]");
System.exit(2);
}
JobConf job = new JobConf();
job.setMapperClass(MapperClass.class);
/** For MapOnly jobs, the number of reducers must be explicitly set to 0. */
job.setNumReduceTasks(0);
/** Configure information about input and output tables. */
InputUtils.addTable(TableInfo.builder().tableName(args[0]).build(), job);
OutputUtils.addTable(TableInfo.builder().tableName(args[1]).build(), job);
if (args.length == 3) {
String options = new String(args[2]);
/** You can specify key-value pairs in the JobConf file, and use getJobConf of the context to query the configurations in a mapper. */
if (options.contains("setup")) {
job.setBoolean("option.mapper.setup", true);
}
if (options.contains("map")) {
job.setBoolean("option.mapper.map", true);
}
if (options.contains("cleanup")) {
job.setBoolean("option.mapper.cleanup", true);
}
}
JobClient.runJob(job);
}
}
Funcionamento do código
main() — configure do job
O uso de job.setNumReduceTasks(0) é obrigatório para jobs map-only. Sem essa definição, o framework espera um reducer e o job falha. Os métodos InputUtils.addTable e OutputUtils.addTable vinculam as tabelas de entrada e saída ao job por meio dos argumentos da linha de comando.
O terceiro argumento opcional (setup, map ou cleanup) defina o sinalizador booleano correspondente no JobConf. Cada método de ciclo de vida lê seu próprio sinalizador via context.getJobConf().getBoolean(...) e grava a saída apenas quando o valor for true.
setup() — executa uma vez antes do início do processamento
Grava um único registro com chave "setup" e contagem 1. A execução ocorre quando option.mapper.setup=true.
map() — executa uma vez por registro de entrada
Lê o primeiro campo de cada registro de entrada (record.get(0)) e o grava com contagem 1. Este método roda quando option.mapper.map=true. Neste exemplo, o argumento map ativa esse método, motivo pelo qual a saída contém duas linhas — uma para cada linha de entrada.
cleanup() — executa uma vez após o processamento de todos os registros
Grava um único registro com chave "cleanup" e contagem 1. A execução acontece quando option.mapper.cleanup=true.