Tous les produits
Search
Centre de documentation

DataWorks:Nœud EMR MR

Dernière mise à jour :Aug 10, 2026

Dans DataWorks, créez un nœud MapReduce (MR) E-MapReduce (EMR) pour diviser les grands ensembles de données en tâches map parallèles, ce qui améliore considérablement l'efficacité du traitement des données. Cette rubrique présente un exemple de développement et de configuration d'un job EMR MR qui lit un fichier texte depuis Object Storage Service (OSS) et en compte les mots.

Prérequis

  • Créez un cluster Alibaba Cloud EMR et enregistrez-le dans DataWorks. Pour plus d'informations, consultez la section New Data Studio : Attacher une ressource de calcul EMR.

  • (Facultatif, pour les utilisateurs RAM) L'utilisateur Resource Access Management (RAM) chargé du développement des tâches doit être ajouté à l'espace de travail et recevoir le rôle Development ou Workspace Administrator (ce rôle inclut des autorisations étendues et doit être attribué avec prudence). Pour plus d'informations, consultez la section Ajouter des membres à un espace de travail.

    Si vous utilisez un compte racine, ignorez cette étape.
  • Pour suivre l'exemple de cette rubrique, créez un bucket dans Object Storage Service (OSS). Pour plus d'informations, consultez la section Créer des buckets.

Limites

  • Exécutez ce type de nœud uniquement sur un groupe de ressources serverless (recommandé) ou un groupe de ressources exclusif pour la planification.

  • Pour gérer les métadonnées d'un cluster DataLake ou personnalisé dans DataWorks, configurez d'abord EMR-HOOK dans le cluster. Pour plus d'informations, consultez la section Configurer Hive EMR-HOOK.

    Remarque

    Si EMR-HOOK n'est pas configuré dans le cluster, DataWorks ne peut pas afficher les métadonnées en temps réel, générer des journaux d'audit, afficher la lignée des données ni effectuer des tâches de gouvernance des données liées à EMR.

Préparer les données initiales et un package JAR

Préparer les données initiales

Créez un exemple de fichier nommé input01.txt avec le contenu suivant.

hadoop emr hadoop dw
hive hadoop
dw emr

Télécharger le fichier de données initiales

  1. Connectez-vous à la console OSS. Dans le volet de navigation de gauche, cliquez sur Buckets.

  2. Cliquez sur le nom du bucket cible pour ouvrir la page File Management.

    Cet exemple utilise un bucket nommé onaliyun-bucket-2.

  3. Cliquez sur Create Directory pour créer les répertoires destinés aux données initiales et à la ressource JAR.

    • Définissez le paramètre Directory Name sur emr/datas/wordcount02/inputs afin de créer le répertoire pour les données initiales.

    • Définissez le paramètre Directory Name sur emr/jars afin de créer le répertoire pour la ressource JAR.

  4. Téléchargez le fichier de données initiales dans son répertoire.

    • Accédez au chemin /emr/datas/wordcount02/inputs et cliquez sur Upload File.

    • Dans la zone Files to Upload, cliquez sur Select Files, ajoutez le fichier input01.txt au bucket, puis cliquez sur Upload File.

Générer un job MapReduce et un package JAR

  1. Ouvrez votre projet IntelliJ IDEA et ajoutez les dépendances suivantes à votre fichier 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. Pour lire et écrire dans un fichier OSS via MapReduce, configurez les paramètres suivants.

    Important

    Avertissement de risque : La paire AccessKey de votre compte Alibaba Cloud accorde un accès complet à toutes les opérations API. La divulgation de votre AccessKey ID et AccessKey Secret compromet la sécurité de toutes les ressources de votre compte. Nous vous recommandons vivement d'utiliser un utilisateur RAM pour les appels API et les opérations quotidiennes. Ne codez pas en dur votre AccessKey ID ou AccessKey Secret dans le code de votre projet ou dans tout autre emplacement accessible publiquement. Le code suivant est fourni à titre de démonstration uniquement. Maintenez la sécurité de vos informations AccessKey.

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

    Le tableau suivant décrit les paramètres.

    • ${accessKeyId} : L'AccessKey ID de votre compte Alibaba Cloud.

    • ${accessKeySecret} : L'AccessKey Secret de votre compte Alibaba Cloud.

    • ${endpoint} : Le point de terminaison public d'OSS. Le point de terminaison dépend de la région où se trouve votre cluster EMR. Le bucket OSS et le cluster doivent résider dans la même région. Pour plus d'informations, consultez la section Régions et points de terminaison.

    Le code Java ci-dessous est une version modifiée de l'exemple officiel Hadoop WordCount. Il inclut les configurations pour l'AccessKey ID et l'AccessKey Secret afin d'autoriser le job à accéder aux fichiers OSS.

    Exemple de code

    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. Après avoir modifié le code Java, empaquetez-le dans un fichier JAR. Cet exemple génère un fichier JAR nommé onaliyun_mr_wordcount-1.0-SNAPSHOT.jar.

Procédure

  1. Dans l'onglet de configuration du nœud EMR MR, développez la tâche comme suit :

    Développer la tâche EMR MR

    Choisissez l'une des méthodes suivantes en fonction de vos besoins :

    Méthode 1 : Télécharger et référencer un JAR

    Vous pouvez également télécharger une ressource depuis votre machine locale vers DataStudio, puis la référencer dans un nœud. Si une ressource est trop volumineuse pour être téléchargée via la console DataWorks, stockez-la dans HDFS et référencez-la dans votre code.

    1. Créez une ressource JAR.

      1. Pour plus d'informations, consultez la section Gérer les ressources. Stockez le package JAR issu de l'étape Préparer les données initiales et un package JAR dans le répertoire emr/jars. Cliquez sur Click Upload.

      2. Configurez les paramètres Storage Path, Data Sources et Resource Group.

      3. Cliquez sur Save.

      image

    2. Référencez la ressource JAR.

      1. Ouvrez le nœud EMR MR pour accéder à son onglet de configuration.

      2. Dans le volet Gestion des ressources à gauche, recherchez la ressource que vous souhaitez référencer. Dans cet exemple, il s'agit de la ressource onaliyun_mr_wordcount-1.0-SNAPSHOT.jar. Faites un clic droit sur la ressource et sélectionnez Insert Resource Path.

      3. Après avoir référencé la ressource, une instruction de référence apparaît dans l'onglet de configuration du nœud EMR MR, indiquant que la ressource a été référencée avec succès. Exécutez la commande suivante. Remplacez le package de ressources, le nom du bucket et le chemin par vos informations réelles.

        ##@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
        Remarque

        Les commentaires ne sont pas pris en charge dans l'éditeur de code des nœuds EMR MR.

    Méthode 2 : Référencer une ressource OSS

    Utilisez la méthode OSS REF pour référencer directement une ressource depuis OSS. Lors de l'exécution du nœud, DataWorks télécharge automatiquement la ressource OSS référencée dans l'environnement local. Cette méthode est souvent utilisée dans les scénarios où une tâche EMR dépend d'un fichier JAR ou d'un script.

    1. Téléchargez la ressource JAR.

      1. Après avoir développé le code, connectez-vous à la console OSS. Dans le volet de navigation de gauche, cliquez sur Buckets.

      2. Cliquez sur le nom du bucket cible pour ouvrir la page File Management.

        Cet exemple utilise un bucket nommé onaliyun-bucket-2.

      3. Téléchargez la ressource JAR dans son répertoire.

        Accédez au répertoire emr/jars. Cliquez sur Upload File. Dans la zone Files to Upload, cliquez sur Select Files, ajoutez le fichier onaliyun_mr_wordcount-1.0-SNAPSHOT.jar, puis cliquez sur Upload File.

    2. Référencez la ressource JAR.

      Dans l'onglet de configuration du nœud EMR MR, écrivez le code permettant de référencer la ressource 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
      Remarque

      La commande respecte le format suivant : hadoop jar <Chemin du JAR à exécuter> <Nom complet de la classe principale> <Répertoire du fichier d'entrée> <Répertoire de sortie>.

      Le tableau suivant décrit les paramètres relatifs au chemin du JAR.

      Paramètre

      Description

      Chemin du JAR à exécuter

      Le format est ossref://{endpoint}/{bucket}/{object}

      • Endpoint : Le point de terminaison public d'OSS. Si ce paramètre est laissé vide, vous ne pouvez référencer des ressources que depuis un bucket situé dans la même région que le cluster EMR.

      • Bucket : Conteneur utilisé par OSS pour stocker les objets. Chaque Bucket possède un nom unique. Vous pouvez vous connecter à la console OSS pour afficher tous les Buckets de votre compte.

      • object : Objet spécifique, qui peut être un fichier ou un chemin, stocké dans un bucket.

    (Facultatif) Configurer les paramètres avancés

    Dans le volet de droite, cliquez sur l'onglet Scheduling Settings. Configurez les paramètres suivants dans la section EMR Node Parameters > DataWorks parameters.

    Remarque
    • Les paramètres avancés disponibles varient selon le type de cluster EMR, comme indiqué dans les tableaux suivants.

    • Configurez des propriétés Spark open source supplémentaires dans l'onglet Scheduling Settings sous la section EMR Node Parameters > Spark parameter.

    Cluster DataLake et personnalisé : EMR sur ECS

    Paramètre

    Description

    queue

    File d'attente dans laquelle les jobs sont soumis. La valeur par défaut est default. Pour plus d'informations sur EMR YARN, consultez la section Configurations de file d'attente de base.

    priority

    Priorité du job. La valeur par défaut est 1.

    FLOW_SKIP_SQL_ANALYZE

    Mode d'exécution des instructions SQL. Valeurs possibles :

    • true : Exécute plusieurs instructions SQL à la fois.

    • false (Par défaut) : Exécute une instruction SQL à la fois.

    Remarque

    Ce paramètre est pris en charge uniquement pour les exécutions de test dans l'environnement de développement des données.

    Autres

    Vous pouvez également ajouter des paramètres de job MR personnalisés dans la section de configuration avancée. Lors de la validation du code, DataWorks ajoute automatiquement les nouveaux paramètres à la commande à l'aide de l'instruction -D key=value.

    Cluster Hadoop : EMR sur ECS

    Paramètre

    Description

    queue

    File d'attente dans laquelle les jobs sont soumis. La valeur par défaut est default. Pour plus d'informations sur EMR YARN, consultez la section Configurations de file d'attente de base.

    priority

    Priorité du job. La valeur par défaut est 1.

    USE_GATEWAY

    Indique si les jobs de ce nœud doivent être soumis via un cluster de passerelle. Valeurs possibles :

    • true : Soumet les jobs via un cluster de passerelle.

    • false (Par défaut) : Ne soumet pas les jobs via un cluster de passerelle. Les jobs sont soumis au nœud maître par défaut.

    Remarque

    Si vous définissez ce paramètre sur true mais que le cluster du nœud n'est pas associé à un cluster de passerelle, la soumission du job EMR échoue.

    Exécuter la tâche

    1. Dans la section Run Configuration, sous Compute Resource, configurez les paramètres Compute Resource et Resource Group.

      Remarque
      • Spécifiez également le nombre d'CUs for Scheduling en fonction des besoins de votre tâche. La valeur par défaut est 0.25.

      • Pour accéder à une source de données via Internet public ou un Virtual Private Cloud (VPC), utilisez un groupe de ressources pour la planification capable de se connecter à la source de données. Pour plus d'informations, consultez la section Solutions de connectivité réseau.

    2. Dans la boîte de dialogue des paramètres de la barre d'outils, sélectionnez la source de données que vous avez créée et cliquez sur Run.

  2. Si vous devez exécuter la tâche de nœud périodiquement, configurez ses propriétés de planification. Pour plus d'informations, consultez la section Configurer les propriétés de planification d'un nœud.

  3. Une fois le nœud configuré, déployez-le. Pour plus d'informations, consultez la section Déployer des nœuds.

  4. Une fois la tâche déployée, consultez son statut dans Operation Center. Pour plus d'informations, consultez la section Présentation d'Operation Center.

Afficher les résultats

  • Connectez-vous à la console OSS. Affichez le fichier de sortie dans le répertoire de destination de votre bucket. Dans cet exemple, le chemin est emr/datas/wordcount02/outputs.目标Bucket

  • Consultez les résultats statistiques dans DataWorks.

    1. Créez un nœud EMR Hive. Pour plus d'informations, consultez la section Créer un nœud pour un flux de travail planifié.

    2. Dans le nœud EMR Hive, créez une table externe Hive mappée aux données d'OSS, puis interrogez les données de la table. Voici un exemple de code :

      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;

      L'image suivante illustre le résultat.运行结果