Esta página documenta as principais classes e métodos do MapReduce no MaxCompute Java SDK (odps-sdk-mapred).
Adicionar a dependência do SDK
Pesquise por odps-sdk-mapred no repositório Maven para encontrar a versão mais recente. Adicione a seguinte dependência ao seu projeto:
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-mapred</artifactId>
<version>0.40.10-public</version>
</dependency>
Tipos de dados
O MapReduce do MaxCompute suporta os tipos de dados listados abaixo. A tabela apresenta o mapeamento para os tipos Java correspondentes.
|
Tipo MaxCompute |
Tipo Java |
|
BIGINT |
Long |
|
STRING |
String |
|
DOUBLE |
Double |
|
BOOLEAN |
Boolean |
|
DATETIME |
Date |
|
DECIMAL |
BigDecimal |
Visão geral das classes
|
Classe |
Descrição |
|
|
Classe base para mappers definidos pelo usuário. Converte registros de tabelas de entrada em pares chave-valor e os encaminha para um reducer. Jobs que ignoram a etapa de reduce e gravam resultados diretamente são chamados de MapOnly jobs. |
|
|
Classe base para reducers definidos pelo usuário. Agrega o conjunto de valores associados a cada chave. |
|
|
Fornece o contexto de execução da tarefa. Passado como parâmetro de entrada para os métodos do ciclo de vida de |
|
|
Envia e gerencie jobs. Suporta envio bloqueante (síncrono) e não bloqueante (assíncrono). |
|
|
Representa uma instância de job em execução. Use-o para acompanhar o status, aguardar a conclusão e recuperar contadores. |
|
|
Armazena a configuração de um job MapReduce. Defina um |
MapperBase
O framework chama os métodos do ciclo de vida do mapper nesta ordem: setup uma vez no início, map uma vez por registro de entrada e cleanup uma vez ao final.
public class WordCountMapper extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
// Initialize resources (for example, load lookup tables)
}
@Override
public void map(long key, Record record, TaskContext context) throws IOException {
// Process each input record and emit key-value pairs.
// MaxCompute uses (long key, Record record), not (KEYIN key, VALUEIN value).
Record mapKey = context.createMapOutputKeyRecord();
Record mapValue = context.createMapOutputValueRecord();
// Populate mapKey and mapValue, then emit:
context.write(mapKey, mapValue);
}
@Override
public void cleanup(TaskContext context) throws IOException {
// Release resources or emit final accumulated results
}
}
|
Método |
Descrição |
|
|
Chamado uma vez antes da primeira chamada a |
|
|
Chamado uma vez por registro de entrada. Emita pares chave-valor com |
|
|
Chamado uma vez após a última chamada a |
ReducerBase
O framework chama os métodos do ciclo de vida do reducer nesta ordem: setup uma vez no início, reduce uma vez por grupo de chaves únicas e cleanup uma vez ao final.
public class WordCountReducer extends ReducerBase {
@Override
public void setup(TaskContext context) throws IOException {
// Initialize resources
}
@Override
public void reduce(Record key, Iterator<Record> values, TaskContext context) throws IOException {
// Aggregate all values for this key
long count = 0;
while (values.hasNext()) {
Record val = values.next();
count += val.getBigint(0);
}
Record output = context.createOutputRecord();
output.set(0, key.getString(0));
output.set(1, count);
context.write(output);
}
@Override
public void cleanup(TaskContext context) throws IOException {
// Release resources
}
}
|
Método |
Descrição |
|
|
Chamado uma vez antes da primeira chamada a |
|
|
Chamado uma vez por grupo de chaves únicas. Todos os valores associados à chave são passados como |
|
|
Chamado uma vez após a última chamada a |
TaskContext
TaskContext é passado para cada método do ciclo de vida e fornece acesso a tabelas de saída, recursos e utilitários do framework.
|
Método |
Descrição |
|
|
Retorna informações sobre as tabelas de saída. |
|
|
Cria um registro para a tabela de saída padrão. |
|
|
Cria um registro para a tabela de saída identificada por |
|
|
Cria um registro para a chave de saída da etapa de map. |
|
|
Cria um registro para o valor de saída da etapa de map. |
|
|
Grava um registro na tabela de saída padrão. Pode ser chamado várias vezes durante a etapa de reduce. |
|
|
Grava um registro na tabela de saída identificada por |
|
|
Emite um par chave-valor durante a etapa de map. Pode ser chamado várias vezes. |
|
|
Lê um recurso de arquivo pelo nome. |
|
|
Lê um recurso de tabela pelo nome. |
|
|
Retorna o contador com o nome especificado. |
|
|
Retorna o contador com o nome especificado no grupo indicado. |
|
|
Envia um heartbeat ao framework MapReduce para evitar timeout do worker. |
Timeout do worker
O timeout padrão do worker é de 10 minutos e não pode ser alterado. Se um worker não chamar progress() em 10 minutos, o framework encerra o worker e a tarefa de map ou reduce falha. Chame progress() periodicamente em tarefas de longa duração para manter os workers ativos. O método envia apenas um heartbeat — não reporta o progresso da tarefa.
JobConf
JobConf armazena toda a configuração de um job MapReduce, incluindo as classes de mapper e reducer, os schemas de chave-valor e as declarações de recursos.
|
Método |
Descrição |
|
|
Define a classe do mapper para o job. |
|
|
Define a classe do reducer para o job. |
|
|
Define uma classe de combiner. O combiner pré-agrega registros com a mesma chave na etapa de map, reduzindo o volume de dados transferidos para os reducers. |
|
|
Define o schema das chaves transmitidas do mapper para o reducer. |
|
|
Define o schema dos valores transmitidos do mapper para o reducer. |
|
|
Especifica as colunas usadas para ordenar as chaves antes de enviá-las aos reducers. |
|
|
Especifica as colunas usadas para agrupar as chaves. As colunas de agrupamento devem ser um subconjunto das colunas de ordenação. |
|
|
Define as colunas de chave de partição. Por padrão, corresponde a todas as colunas de chave. |
|
|
Declara os recursos disponíveis para mappers e reducers. Um mapper ou reducer só pode ler recursos declarados aqui. |
|
|
Define o tamanho do split de entrada em MB. Padrão: 256 MB. |
|
|
Define o número de tarefas de reduce. Padrão: um quarto do número de tarefas de map. |
|
|
Define a memória por worker de map em MB. Padrão: 2048 MB. |
|
|
Define a memória por worker de reduce em MB. Padrão: 2048 MB. |
Como as etapas de map e reduce distribuem e agrupam registros
Na etapa de map, o framework calcula um hash de cada registro de saída com base nas colunas de chave de partição para determinar qual reducer o receberá. Os registros são ordenados pelas colunas de ordenação antes do envio.
Na etapa de reduce, os registros são agrupados pelas colunas de agrupamento. Todos os registros que compartilham a mesma chave de agrupamento são passados juntos para uma única chamada a reduce().
As colunas de agrupamento são selecionadas a partir das colunas de ordenação. As colunas de ordenação e as colunas de chave de partição devem existir nas chaves.
JobClient
|
Método |
Descrição |
|
|
Envia um job em modo bloqueante (síncrono). Aguarda até a conclusão do job e retorna um |
|
|
Envia um job em modo não bloqueante (assíncrono). Retorna imediatamente com um |
RunningJob
|
Método |
Descrição |
|
|
Retorna o ID da instância do job. Use esse ID para visualize logs operacionais e gerencie o job. |
|
|
Retorna |
|
|
Retorna |
|
|
Aguarda o término de uma instância de job. O método é utilizado para jobs enviados em modo síncrono. |
|
|
Retorna o status atual da instância do job. |
|
|
Encerra o job em execução. |
|
|
Retorna todos os dados de contadores do job. |
InputUtils
|
Método |
Descrição |
|
|
Adiciona uma única tabela de entrada ao job. Chamadas sucessivas acrescentam cada tabela à fila de entrada. |
|
|
Define várias tabelas de entrada de uma só vez. |
OutputUtils
|
Método |
Descrição |
|
|
Adiciona uma única tabela de saída ao job. Chamadas sucessivas acrescentam cada tabela à fila de saída. |
|
|
Define várias tabelas de saída de uma só vez. |
Pipeline (modelo MapReduce estendido)
Pipeline é o ponto de entrada para o modelo MapReduce estendido, que permite encadear múltiplos mappers e reducers em um único job. Use Pipeline.builder() para construir o pipeline.
Métodos do Builder
public Builder addMapper(Class<? extends Mapper> mapper)
public Builder addMapper(Class<? extends Mapper> mapper,
Column[] keySchema, Column[] valueSchema, String[] sortCols,
SortOrder[] order, String[] partCols,
Class<? extends Partitioner> theClass, String[] groupCols)
public Builder addReducer(Class<? extends Reducer> reducer)
public Builder addReducer(Class<? extends Reducer> reducer,
Column[] keySchema, Column[] valueSchema, String[] sortCols,
SortOrder[] order, String[] partCols,
Class<? extends Partitioner> theClass, String[] groupCols)
public Builder setOutputKeySchema(Column[] keySchema)
public Builder setOutputValueSchema(Column[] valueSchema)
public Builder setOutputKeySortColumns(String[] sortCols)
public Builder setOutputKeySortOrder(SortOrder[] order)
public Builder setPartitionColumns(String[] partCols)
public Builder setPartitionerClass(Class<? extends Partitioner> theClass)
public Builder setOutputGroupingColumns(String[] cols)
Exemplo
O exemplo a seguir encadeia um mapper e dois reducers (TokenizerMapper → SumReducer → IdentityReducer):
Job job = new Job();
Pipeline pipeline = Pipeline.builder()
.addMapper(TokenizerMapper.class)
.setOutputKeySchema(
new Column[] { new Column("word", OdpsType.STRING) })
.setOutputValueSchema(
new Column[] { new Column("count", OdpsType.BIGINT) })
.addReducer(SumReducer.class)
.setOutputKeySchema(
new Column[] { new Column("count", OdpsType.BIGINT) })
.setOutputValueSchema(
new Column[] { new Column("word", OdpsType.STRING),
new Column("count", OdpsType.BIGINT) })
.addReducer(IdentityReducer.class).createPipeline();
job.setPipeline(pipeline);
job.addInput(...)
job.addOutput(...)
job.submit();
Para encadear um único mapper com um reducer, useJobConfem vez dePipeline. Familiarize-se com a API MapReduce padrão antes de utilizar o modelo estendido.