Todos os produtos
Search
Central de documentação

DataWorks:Nó EMR MR

Última atualização: Jun 27, 2026

No DataWorks, é possível criar um nó E-MapReduce (EMR) MapReduce (MR) para dividir grandes conjuntos de dados em tarefas de mapa paralelas, o que melhora significativamente a eficiência do processamento de dados. Este tópico apresenta um exemplo de desenvolvimento e configuração de um job EMR MR que lê um arquivo de texto do Object Storage Service (OSS) e conta as palavras contidas nele.

Pré-requisitos

  • Crie um cluster Alibaba Cloud EMR e registre-o no DataWorks. Para mais informações, consulte New Data Studio: Attach an EMR compute resource.

  • (Opcional, para usuários RAM) O usuário do Resource Access Management (RAM) responsável pelo desenvolvimento da tarefa deve ser adicionado ao workspace e receber a função Development ou Workspace Administrator (esta função inclui permissões amplas e deve ser concedida com cautela). Para mais informações, consulte Add workspace members.

    Se você estiver usando uma conta raiz, pule esta etapa.
  • Para acompanhar o exemplo deste tópico, crie um bucket no Object Storage Service (OSS). Para mais informações, consulte Create buckets.

Limitações

  • A execução deste tipo de nó é suportada apenas em um serverless resource group (recomendado) ou em um grupo de recursos exclusivo para agendamento.

  • Para gerenciar metadados de um cluster DataLake ou personalizado no DataWorks, configure primeiro o EMR-HOOK no cluster. Para mais informações, consulte Configure Hive EMR-HOOK.

    Nota

    Sem a configuração do EMR-HOOK no cluster, o DataWorks não consegue exibir metadados em tempo real, gerar logs de auditoria, mostrar linhagem de dados ou executar tarefas de governança de dados relacionadas ao EMR.

Preparar dados iniciais e pacote JAR

Preparar os dados iniciais

Crie um arquivo de amostra chamado input01.txt com o seguinte conteúdo.

hadoop emr hadoop dw
hive hadoop
dw emr

Fazer upload do arquivo de dados iniciais

  1. Faça login no OSS console. No painel de navegação à esquerda, clique em Buckets.

  2. Clique no nome do bucket desejado para abrir a página File Management.

    Este exemplo utiliza um bucket chamado onaliyun-bucket-2.

  3. Clique em Create Directory para criar diretórios destinados aos dados iniciais e ao recurso JAR.

    • Defina Directory Name como emr/datas/wordcount02/inputs para criar o diretório dos dados iniciais.

    • Defina Directory Name como emr/jars para criar o diretório do recurso JAR.

  4. Faça o upload do arquivo de dados iniciais para seu respectivo diretório.

    • Acesse o caminho /emr/datas/wordcount02/inputs e clique em Upload File.

    • Na área Files to Upload, clique em Select Files, adicione o arquivo input01.txt ao bucket e, em seguida, clique em Upload File.

Criar job MapReduce e pacote JAR

  1. Abra seu projeto no IntelliJ IDEA e adicione as seguintes dependências ao arquivo pom.xml.

            <dependency>
                <groupId>org.apache.hadoop</groupId>
                <artifactId>hadoop-mapreduce-client-common</artifactId>
                <version>2.8.5</version> <!--Use version 2.8.5, which is the version used by EMR MR.-->
            </dependency>
            <dependency>
                <groupId>org.apache.hadoop</groupId>
                <artifactId>hadoop-common</artifactId>
                <version>2.8.5</version>
            </dependency>
  2. Para ler e gravar em arquivos do OSS via MapReduce, configure os seguintes parâmetros.

    Importante

    Aviso de risco: O par de AccessKey da sua conta Alibaba Cloud concede acesso total a todas as operações de API. A exposição do seu AccessKey ID e AccessKey Secret compromete a segurança de todos os recursos sob sua conta. Recomendamos fortemente o uso de um usuário RAM para chamadas de API e operações diárias. Não codifique rigidamente (hardcode) seu AccessKey ID ou AccessKey Secret no código do projeto ou em qualquer outro local publicamente acessível. O código abaixo serve apenas para fins de demonstração. Mantenha suas informações de AccessKey seguras.

    conf.set("fs.oss.accessKeyId", "${accessKeyId}");
    conf.set("fs.oss.accessKeySecret", "${accessKeySecret}");
    conf.set("fs.oss.endpoint","${endpoint}");

    A tabela a seguir descreve os parâmetros.

    • ${accessKeyId}: O AccessKey ID da sua conta Alibaba Cloud.

    • ${accessKeySecret}: O AccessKey Secret da sua conta Alibaba Cloud.

    • ${endpoint}: O endpoint público do OSS. O endpoint depende da região onde seu cluster EMR está localizado. O bucket do OSS e o cluster devem estar na mesma região. Para mais informações, consulte Regions and endpoints.

    O código Java a seguir é uma versão modificada do exemplo oficial Hadoop WordCount. Ele inclui configurações para o AccessKey ID e AccessKey Secret, concedendo ao job permissão para acessar arquivos no OSS.

    Código de exemplo

    package cn.apache.hadoop.onaliyun.examples;
    
    import java.io.IOException;
    import java.util.StringTokenizer;
    
    import org.apache.hadoop.conf.Configuration;
    import org.apache.hadoop.fs.Path;
    import org.apache.hadoop.io.IntWritable;
    import org.apache.hadoop.io.Text;
    import org.apache.hadoop.mapreduce.Job;
    import org.apache.hadoop.mapreduce.Mapper;
    import org.apache.hadoop.mapreduce.Reducer;
    import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
    import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
    import org.apache.hadoop.util.GenericOptionsParser;
    
    public class EmrWordCount {
        public static class TokenizerMapper
                extends Mapper<Object, Text, Text, IntWritable> {
            private final static IntWritable one = new IntWritable(1);
            private Text word = new Text();
    
            public void map(Object key, Text value, Context context
            ) throws IOException, InterruptedException {
                StringTokenizer itr = new StringTokenizer(value.toString());
                while (itr.hasMoreTokens()) {
                    word.set(itr.nextToken());
                    context.write(word, one);
                }
            }
        }
    
        public static class IntSumReducer
                extends Reducer<Text, IntWritable, Text, IntWritable> {
            private IntWritable result = new IntWritable();
    
            public void reduce(Text key, Iterable<IntWritable> values,
                               Context context
            ) throws IOException, InterruptedException {
                int sum = 0;
                for (IntWritable val : values) {
                    sum += val.get();
                }
                result.set(sum);
                context.write(key, result);
            }
        }
    
        public static void main(String[] args) throws Exception {
            Configuration conf = new Configuration();
            String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs();
            if (otherArgs.length < 2) {
                System.err.println("Usage: wordcount <in> [<in>...] <out>");
                System.exit(2);
            }
            conf.set("fs.oss.accessKeyId", "${accessKeyId}"); // 
            conf.set("fs.oss.accessKeySecret", "${accessKeySecret}"); // 
            conf.set("fs.oss.endpoint", "${endpoint}"); //
            Job job = Job.getInstance(conf, "word count");
            job.setJarByClass(EmrWordCount.class);
            job.setMapperClass(TokenizerMapper.class);
            job.setCombinerClass(IntSumReducer.class);
            job.setReducerClass(IntSumReducer.class);
            job.setOutputKeyClass(Text.class);
            job.setOutputValueClass(IntWritable.class);
            for (int i = 0; i < otherArgs.length - 1; ++i) {
                FileInputFormat.addInputPath(job, new Path(otherArgs[i]));
            }
            FileOutputFormat.setOutputPath(job,
                    new Path(otherArgs[otherArgs.length - 1]));
            System.exit(job.waitForCompletion(true) ? 0 : 1);
        }
    }
                                    
  3. Após editar o código Java, empacote-o em um arquivo JAR. Este exemplo gera um arquivo JAR chamado onaliyun_mr_wordcount-1.0-SNAPSHOT.jar.

Procedimento

  1. Na aba de configuração do nó EMR MR, desenvolva a tarefa conforme descrito abaixo:

    Desenvolver a tarefa EMR MR

    Escolha um dos métodos a seguir conforme sua necessidade:

    Método 1: Fazer upload e referenciar JAR

    Também é possível fazer upload de um recurso da máquina local para o DataStudio e referenciá-lo em um nó. Caso o recurso seja muito grande para upload via console DataWorks, armazene-o no HDFS e referencie-o no código.

    1. Crie um recurso JAR.

      1. Para mais informações, consulte Manage resources. Armazene o pacote JAR da etapa Preparar dados iniciais e pacote JAR no diretório emr/jars. Clique em Click Upload.

      2. Configure os parâmetros Storage Path, Data Sources e Resource Group.

      3. Clique em Save.

      image

    2. Referencie o recurso JAR.

      1. Abra o nó EMR MR para acessar sua aba de configuração.

      2. No painel Resource Management à esquerda, localize o recurso que deseja referenciar. Neste exemplo, o recurso é onaliyun_mr_wordcount-1.0-SNAPSHOT.jar. Clique com o botão direito no recurso e selecione Insert Resource Path.

      3. Após referenciar o recurso, uma instrução de referência aparecerá na aba de configuração do nó EMR MR, indicando sucesso na operação. Execute o comando abaixo, substituindo o pacote de recursos, o nome do bucket e o caminho pelas suas informações reais.

        ##@resource_reference{"onaliyun_mr_wordcount-1.0-SNAPSHOT.jar"}
        onaliyun_mr_wordcount-1.0-SNAPSHOT.jar cn.apache.hadoop.onaliyun.examples.EmrWordCount oss://onaliyun-bucket-2/emr/datas/wordcount02/inputs oss://onaliyun-bucket-2/emr/datas/wordcount02/outputs
        Nota

        O editor de código para nós EMR MR não suporta comentários.

    Método 2: Referenciar recurso OSS

    Utilize o método OSS REF para referenciar diretamente um recurso do OSS. Durante a execução do nó, o DataWorks baixa automaticamente o recurso OSS referenciado para o ambiente local. Essa abordagem é comum em cenários onde uma tarefa EMR depende de um arquivo JAR ou script.

    1. Faça o upload do recurso JAR.

      1. Após desenvolver o código, faça login no OSS console. No painel de navegação à esquerda, clique em Buckets.

      2. Clique no nome do bucket desejado para abrir a página File Management.

        Este exemplo utiliza um bucket chamado onaliyun-bucket-2.

      3. Faça o upload do recurso JAR para seu diretório.

        Acesse o diretório emr/jars. Clique em Upload File. Na área Files to Upload, clique em Select Files, adicione o arquivo onaliyun_mr_wordcount-1.0-SNAPSHOT.jar e, em seguida, clique em Upload File.

    2. Referencie o recurso JAR.

      Na aba de configuração do nó EMR MR, escreva o código para referenciar o recurso JAR.

      hadoop jar ossref://onaliyun-bucket-2/emr/jars/onaliyun_mr_wordcount-1.0-SNAPSHOT.jar cn.apache.hadoop.onaliyun.examples.EmrWordCount oss://onaliyun-bucket-2/emr/datas/wordcount02/inputs oss://onaliyun-bucket-2/emr/datas/wordcount02/outputs
      Nota

      O comando segue o formato: hadoop jar <Caminho do JAR a ser executado> <Nome completo da classe principal> <Diretório do arquivo de entrada> <Diretório de saída>.

      A tabela a seguir descreve os parâmetros do caminho do JAR.

      Parâmetro

      Descrição

      Caminho do JAR a ser executado

      O formato é ossref://{endpoint}/{bucket}/{object}

      • Endpoint: O endpoint público do OSS. Se este parâmetro for deixado em branco, só será possível referenciar recursos de um bucket na mesma região do cluster EMR.

      • Bucket: Um contêiner usado pelo OSS para armazenar objetos. Cada Bucket possui um nome único. Faça login no OSS console para visualizar todos os Buckets da sua conta.

      • object: Um objeto específico, que pode ser um arquivo ou um caminho, armazenado em um bucket.

    (Opcional) Configurar parâmetros avançados

    No painel à direita, clique na aba Scheduling Settings. Configure os seguintes parâmetros na seção EMR Node Parameters > DataWorks parameters.

    Nota
    • Os parâmetros avançados disponíveis variam conforme o tipo de cluster EMR, conforme mostrado nas tabelas a seguir.

    • Configure propriedades adicionais do Spark open-source na aba Scheduling Settings, dentro da seção EMR Node Parameters > Spark parameter.

    Cluster DataLake e personalizado: EMR on ECS

    Parâmetro

    Descrição

    queue

    A fila onde os jobs são submetidos. O valor padrão é default. Para mais informações sobre o EMR YARN, consulte Basic queue configurations.

    priority

    A prioridade do job. O valor padrão é 1.

    FLOW_SKIP_SQL_ANALYZE

    O modo de execução para instruções SQL. Valores válidos:

    • true: Executa múltiplas instruções SQL simultaneamente.

    • false (Padrão): Executa uma instrução SQL por vez.

    Nota

    Este parâmetro é suportado apenas para execuções de teste no ambiente de desenvolvimento de dados.

    Outros

    Também é possível adicionar parâmetros personalizados de job MR na seção de configuração avançada. Ao enviar o código, o DataWorks adiciona automaticamente os novos parâmetros ao comando usando a instrução -D key=value.

    Cluster Hadoop: EMR on ECS

    Parâmetro

    Descrição

    queue

    A fila onde os jobs são submetidos. O valor padrão é default. Para mais informações sobre o EMR YARN, consulte Basic queue configurations.

    priority

    A prioridade do job. O valor padrão é 1.

    USE_GATEWAY

    Define se os jobs deste nó serão submetidos através de um cluster gateway. Valores válidos:

    • true: Submete jobs através de um cluster gateway.

    • false (Padrão): Não submete jobs através de um cluster gateway. Por padrão, os jobs são submetidos ao nó mestre.

    Nota

    Se você definir este parâmetro como true, mas o cluster do nó não estiver associado a um cluster gateway, a submissão do job EMR falhará.

    Executar a tarefa

    1. Em Run Configuration, dentro de Compute Resource, configure Compute Resource e Resource Group.

      Nota
      • Especifique as CUs for Scheduling conforme as necessidades da sua tarefa. O valor padrão é 0.25.

      • Para acessar uma fonte de dados pela internet pública ou em uma Virtual Private Cloud (VPC), utilize um grupo de recursos para agendamento que tenha conectividade com essa fonte. Para mais informações, consulte Network connectivity solutions.

    2. Na caixa de diálogo de parâmetros da barra de ferramentas, selecione a fonte de dados criada e clique em Run.

  2. Caso precise executar a tarefa do nó periodicamente, configure suas propriedades de agendamento. Para mais informações, consulte Configure scheduling properties for a node.

  3. Após configurar o nó, faça o deploy dele. Para mais informações, consulte Deploy nodes.

  4. Depois que a tarefa for implantada, visualize seu status no Operation Center. Para mais informações, consulte Introduction to Operation Center.

Visualizar os resultados

  • Faça login no OSS console. O arquivo de saída estará disponível no diretório de destino dentro do seu bucket. Neste exemplo, o caminho é emr/datas/wordcount02/outputs.目标Bucket

  • Leia os resultados estatísticos no DataWorks.

    1. Crie um nó EMR Hive. Para mais informações, consulte Create a node for a scheduled workflow.

    2. No nó EMR Hive, crie uma tabela externa do Hive mapeada para os dados no OSS e consulte os dados da tabela. Veja abaixo um exemplo de código:

      CREATE EXTERNAL TABLE IF NOT EXISTS wordcount02_result_tb
      (
          `word` STRING COMMENT 'Word',
          `count` STRING COMMENT 'Count'   
      ) 
      ROW FORMAT delimited fields terminated by '\t'
      location 'oss://onaliyun-bucket-2/emr/datas/wordcount02/outputs/';
      
      SELECT * FROM wordcount02_result_tb;

      A figura a seguir mostra o resultado.运行结果