Este exemplo demonstra como encadear vários jobs MapReduce sequencialmente no MaxCompute. Como um único job MapReduce não pode invocar outro durante a execução, caso sua lógica de processamento exija iterações em que cada etapa dependa do resultado da anterior, orquestre múltiplos jobs a partir de um método main no lado do cliente. Este exemplo usa um contador para controlar o loop de iteração.
Funcionamento
O programa executa dois tipos de jobs em sequência:
Job de inicialização (
InitMapper): grava um valor inicial de contador (2) na tabela de saída.Jobs de decremento (
DecreaseMapper): executam em loop. A cada iteração, o sistema lê a saída anterior de uma tabela de recursos, subtrai 1 do valor e defina um contador. O loop termina quando o contador atinge0.
A transferência de dados entre jobs utiliza dois mecanismos:
Tabela de recursos (
multijobs_res_table): transporta o valor numérico entre os jobs.Contador (
multijobs.value): sinaliza ao métodomainse as iterações devem continuar.
O método main atua como único orquestrador: envia os jobs em ordem e verifica a conclusão de cada etapa antes de prosseguir.
Pré-requisitos
Antes de começar, verifique se você:
Concluiu a configuração do ambiente descrita em Primeiros passos.
Compilou o pacote
mapreduce-examples.jare o salvou embin\data\resourcesno diretório de instalação do cliente MaxCompute.
Preparar tabelas e recursos
-
Crie as tabelas de teste:
CREATE TABLE mr_empty (key STRING, value STRING); CREATE TABLE mr_multijobs_out (value BIGINT); -
Registre os recursos utilizados pelo job:
add table mr_multijobs_out as multijobs_res_table -f; -- Omit -f when adding the JAR for the first time. add jar data\resources\mapreduce-examples.jar -f;A tabela
mr_multijobs_outé registrada comomultijobs_res_tablepara permitir que oDecreaseMapperleia a saída do job anterior como uma tabela de recursos. O registro do JAR permite que o cliente MaxCompute localize a classe durante a execução.
Executar MultiJobs
Execute o seguinte comando no cliente MaxCompute:
jar -resources mapreduce-examples.jar,multijobs_res_table -classpath data\resources\mapreduce-examples.jar \
com.aliyun.odps.mapred.open.example.MultiJobs mr_multijobs_out;
Parâmetros do comando:
|
Parâmetro |
Valor |
Descrição |
|
|
|
Recursos disponíveis para o job: o pacote JAR e a tabela de recursos com a saída do job anterior. |
|
|
|
Caminho local do JAR que o cliente usa para localizar a classe principal. |
|
Classe principal |
|
Ponto de entrada do programa. |
|
Argumento |
|
Nome da tabela de saída passado ao método |
Resultado esperado
Após a conclusão do job, a tabela mr_multijobs_out contém um registro:
+------------+
| value |
+------------+
| 0 |
+------------+
Código de exemplo
Para configurar as dependências do Project Object Model (POM), consulte a seção Precauções no guia Primeiros passos.
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.data.TableInfo;
import com.aliyun.odps.mapred.JobClient;
import com.aliyun.odps.mapred.MapperBase;
import com.aliyun.odps.mapred.RunningJob;
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;
/**
* MultiJobs
*
* Running multiple job
*
**/
public class MultiJobs {
public static class InitMapper extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
Record record = context.createOutputRecord();
long v = context.getJobConf().getLong("multijobs.value", 2);
record.set(0, v);
context.write(record);
}
}
public static class DecreaseMapper extends MapperBase {
@Override
public void cleanup(TaskContext context) throws IOException {
/** Obtain the variable values that are defined in the main function from JobConf. */
long expect = context.getJobConf().getLong("multijobs.expect.value", -1);
long v = -1;
int count = 0;
/** Read the data from the output table of the previous job. */
Iterator<Record> iter = context.readResourceTable("multijobs_res_table");
while (iter.hasNext()) {
Record r = iter.next();
v = (Long) r.get(0);
if (expect != v) {
throw new IOException("expect: " + expect + ", but: " + v);
}
count++;
}
if (count != 1) {
throw new IOException("res_table should have 1 record, but: " + count);
}
Record record = context.createOutputRecord();
v--;
record.set(0, v);
context.write(record);
/** Set the counter. The counter value can be obtained in the main function after the job is completed. */
context.getCounter("multijobs", "value").setValue(v);
}
}
public static void main(String[] args) throws Exception {
if (args.length != 1) {
System.err.println("Usage: TestMultiJobs <table>");
System.exit(1);
}
String tbl = args[0];
long iterCount = 2;
System.err.println("Start to run init job.");
JobConf initJob = new JobConf();
initJob.setLong("multijobs.value", iterCount);
initJob.setMapperClass(InitMapper.class);
InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), initJob);
OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), initJob);
initJob.setMapOutputKeySchema(SchemaUtils.fromString("key:string"));
initJob.setMapOutputValueSchema(SchemaUtils.fromString("value:string"));
/** Explicitly set the number of reducers to 0 for map-only jobs. */
initJob.setNumReduceTasks(0);
JobClient.runJob(initJob);
while (true) {
System.err.println("Start to run iter job, count: " + iterCount);
JobConf decJob = new JobConf();
decJob.setLong("multijobs.expect.value", iterCount);
decJob.setMapperClass(DecreaseMapper.class);
InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), decJob);
OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), decJob);
/** Explicitly set the number of reducers to 0 for map-only jobs. */
decJob.setNumReduceTasks(0);
RunningJob rJob = JobClient.runJob(decJob);
iterCount--;
/** If the specified number of iterations is reached, exit the loop. */
if (rJob.getCounters().findCounter("multijobs", "value").getValue() == 0) {
break;
}
}
if (iterCount != 0) {
throw new IOException("Job failed.");
}
}
}
Explicação do código
InitMapper
Executa uma única vez no método setup(). Lê o valor inicial do contador (multijobs.value, padrão 2) do JobConf e o grava na tabela de saída. Trata-se de um job apenas de mapa; a chamada setNumReduceTasks(0) impede que o framework inicie qualquer reducer.
DecreaseMapper
Executa no método cleanup() ao final de cada iteração. Este componente realiza as seguintes ações:
Lê o valor esperado (
multijobs.expect.value) doJobConf.Recupera o único registro da tabela
multijobs_res_tableusandocontext.readResourceTable(). Lança uma exceçãoIOExceptionse a quantidade de registros não for exatamente 1 ou se o valor divergir do esperado.Decrementa o valor em 1, grava o resultado na tabela de saída e atualiza o contador (
multijobs.value) com o novo valor.
Método main**
Controla a sequência de jobs:
Envia o job de inicialização e aguarda sua conclusão.
Inicia um loop que submete os jobs de decremento individualmente e espera o término de cada um antes de prosseguir.
Verifica o contador
multijobs.valueapós cada job de decremento. Se o valor for igual a0, interrompe o loop.Valida se
iterCounté0ao final do loop. Caso contrário, lança uma exceçãoIOExceptionindicando falha no job.