Tous les produits
Search
Centre de documentation

E-MapReduce:Basic Usage

Dernière mise à jour :Aug 09, 2026

Découvrez comment soumettre et afficher des jobs Flink sur E-MapReduce.

Contexte

Le service Flink d'un cluster Dataflow est déployé en mode YARN. Connectez-vous au cluster Dataflow via SSH pour soumettre des jobs Flink depuis la ligne de commande.

Un cluster Dataflow déployé en mode YARN prend en charge la soumission de jobs Flink en mode session, en mode cluster par job et en mode application.

Mode

Description

Avantages et inconvénients

Processus général

La figure suivante illustre le processus général de soumission et d'affichage d'un job Flink.

Par exemple, si une exception dans un job provoque l'arrêt d'un TaskManager, tous les autres jobs s'exécutant sur ce TaskManager échoueront. De plus, comme un cluster ne dispose que d'un seul JobManager, la charge sur le JobManager augmente avec le nombre de jobs.

  • Avantage : Lors de la soumission des jobs, la surcharge temporelle liée à l'allocation des ressources est moindre que dans les autres modes.

  • Inconvénient : Comme tous les jobs s'exécutent dans le même cluster, ils entrent en concurrence pour les ressources et peuvent s'interférer mutuellement.

Compte tenu de ces caractéristiques, ce modèle convient au déploiement de jobs dont le temps de démarrage et la durée d'exécution sont relativement courts.

Mode cluster par job

En mode cluster par job, chaque fois qu'un job Flink est soumis, YARN démarre un nouveau cluster Flink, puis exécute le job. Lorsque le job termine son exécution ou est annulé, le cluster Flink est également libéré.

  • Avantages : Isolation des ressources entre les jobs ; le comportement anormal d'un job n'affecte pas les autres jobs.

    Chaque job correspondant à un JobManager unique, il n'y a pas de risque qu'un JobManager exécute plusieurs jobs et génère des problèmes de charge élevée.

  • Inconvénients : Chaque exécution de job nécessite le démarrage d'un cluster Flink dédié, ce qui entraîne une surcharge de démarrage plus importante.

Sur la base des caractéristiques ci-dessus, ce mode convient généralement aux jobs de longue durée.

Mode application

En mode application, chaque fois que vous soumettez une application Flink (une application contient un ou plusieurs jobs), YARN démarre un nouveau cluster Flink. Lorsque l'application termine son exécution ou est annulée, le cluster Flink de l'application est également libéré.

Ce mode diffère du mode par job en ce sens que la méthode main() du fichier JAR correspondant à l'application est exécutée sur le JobManager dans le cluster.

Si le fichier JAR soumis contient plusieurs jobs, tous ces jobs s'exécuteront dans le cluster de l'application.

  • Avantages : permet de réduire la charge sur le client lors de la soumission des jobs.

  • Inconvénient : Chaque exécution d'une application Flink nécessite le démarrage d'un cluster Flink dédié, ce qui prend plus de temps.

Prérequis

Un cluster Dataflow a été créé en mode Flink. Pour plus d'informations, consultez Créer un cluster.

Soumettre et afficher des jobs Flink

Remarque

Cette rubrique utilise l'exemple Flink TopSpeedWindowing. Il s'agit d'un job de streaming de longue durée.

Vous pouvez choisir parmi les trois modes suivants pour soumettre et afficher les jobs :

Mode session

  1. Connectez-vous au nœud maître du cluster via SSH. Pour plus d'informations, consultez Se connecter au nœud maître d'un cluster.

  2. Exécutez la commande suivante pour démarrer une session YARN.

    yarn-session.sh --detached

    Une fois la commande exécutée avec succès, le système renvoie un ID d'application. Par exemple, application_1750137174986_0001. Cet ID est désigné par <application_XXXX_YY> dans les sections suivantes.

    mr.aliyuncs.com:33879 of application 'application_1750137174986_0001'.
    JobManager Web Interface: http://core-1-1.c-1f6ec9xxx.cn-hangzhou.emr.aliyuncs.com:33879
    2025-06-17 13:19:20,152 INFO  org.apache.flink.yarn.cli.FlinkYarnSessionCli                [] - The Flink YARN session cluster has been started in detached mode. In order to stop Flink gracefully, use the following command:
    $ echo "stop" | ./bin/yarn-session.sh -id application_1750137174986_0001
    If this should not be possible, then you can also kill Flink via YARN's web interface or via:
    $ yarn application -kill application_1750137174986_0001
    Note that killing Flink might not clean up all job artifacts and temporary files.
  3. Exécutez la commande suivante pour soumettre le job.

    flink run --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar

    Après la soumission du job, le système renvoie un message similaire au suivant.

    [root@master-1-1(172.17.xxx.xxx) ~]# flink run --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar
    SLF4J: Class path contains multiple SLF4J bindings.
    SLF4J: Found binding in [jar:file:/opt/apps/FLINK/flink-1.17.2-1.0.10/lib/log4j-slf4j-impl-2.17.1.jar!/org/slf4j/impl/StaticLoggerBinder.class]
    SLF4J: Found binding in [jar:file:/opt/apps/HADOOP-COMMON/hadoop-3.2.1-1.3.2-alinux3/share/hadoop/common/lib/slf4j-log4j12-1.7.25.jar!/org/slf4j/impl/StaticLoggerBinder.class]
    SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation.
    SLF4J: Actual binding is of type [org.apache.logging.slf4j.Log4jLoggerFactory]
    2025-06-17 13:29:00,205 INFO  org.apache.flink.yarn.cli.FlinkYarnSessionCli                [] - Found Yarn properties file under /tmp/.yarn-properties-root.
    2025-06-17 13:29:00,205 INFO  org.apache.flink.yarn.cli.FlinkYarnSessionCli                [] - Found Yarn properties file under /tmp/.yarn-properties-root.
    Executing example with default input data.
    Use --input to specify file input.
    Printing result to stdout. Use --output to specify output path.
    2025-06-17 13:29:00,667 WARN  org.apache.flink.yarn.configuration.YarnLogConfigUtil        [] - The configuration directory ('/etc/taihao-apps/flink-conf') already contains a LOG4J config file.If you want to use logback, then please delete or rename the log configuration file.
    2025-06-17 13:29:00,864 INFO  org.apache.hadoop.yarn.client.RMProxy                        [] - Connecting to ResourceManager at master-1-1.c-1f6ec9192d1528ec.cn-hangzhou.emr.aliyuncs.com/172.17.xxx.xxx:8032
    2025-06-17 13:29:01,061 INFO  org.apache.hadoop.yarn.client.AHSProxy                       [] - Connecting to Application History server at master-1-1.c-1f6ec9192d1528ec.cn-hangzhou.emr.aliyuncs.com/172.17.xxx.xxx:10200
    2025-06-17 13:29:01,072 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - No path for the flink jar passed. Using the location of class org.apache.flink.yarn.YarnClusterDescriptor to locate the jar
    2025-06-17 13:29:01,208 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Found Web Interface core-1-1.c-1f6ecxxx.cn-hangzhou.emr.aliyuncs.com:33879 of application 'application_1750137174986_0001'.
    Job has been submitted with JobID 3785db18d371326758d7843dd2a1xxx

    Dans le message, 3785db18d371326758d7843dd2a1**** correspond à l'ID du job. Cet ID est désigné par <jobId> dans les sections suivantes.

  4. Exécutez la commande suivante pour afficher l'état du job.

    flink list -t yarn-session -Dyarn.application.id=<application_XXXX_YY>

    Un message similaire au suivant est renvoyé.

    ------------------ Running/Restarting Jobs -------------------
    16.06.2025 18:20:55 : 3785db18d371326758d7843dd2a1**** : CarTopSpeedWindowingExample (RUNNING)

    Vous pouvez également consulter l'état du job sur l'interface web. Pour plus d'informations, consultez Afficher l'état du job sur l'interface web.

  5. Exécutez la commande suivante pour arrêter le job.

    flink cancel -t yarn-session -Dyarn.application.id=<application_XXXX_YY> <jobId>

Mode cluster par job

  1. Connectez-vous au nœud maître du cluster via SSH. Pour plus d'informations, consultez Se connecter au nœud maître d'un cluster.

  2. Exécutez la commande suivante pour soumettre le job.

    flink run -t yarn-per-job --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar

    Après la soumission du job, le système renvoie un message similaire au suivant.

    $ yarn application -kill application_1750125819948_0003
    Note that killing Flink might not clean up all job artifacts and temporary files.
    2025-06-17 10:44:46,268 INFO  org.apache.flink.yarn.YarnClusterDescriptor                 [] - Found Web Interface core-1-1.c-b9693c.xxx.cn-hangzhou.emr.aliyuncs.com:38037 of application 'application_1750125819948_0003'.
    Job has been submitted with JobID 451aded93de19d6cd238ed3b466xxx
    You have new mail in /var/spool/mail/root

    Dans le message, application_1750125819948_**** correspond à l'ID d'application, désigné par <application_XXXX_YY> dans les sections suivantes. f5f980ac631192b02548235f1bbe**** correspond à l'ID du job, désigné par <jobId> dans les sections suivantes.

  3. Exécutez la commande suivante pour afficher l'état du job.

    flink list -t yarn-per-job -Dyarn.application.id=<application_XXXX_YY>

    Vous pouvez également consulter l'état du job sur l'interface web. Pour plus d'informations, consultez Afficher l'état du job sur l'interface web.

  4. Exécutez la commande suivante pour arrêter le job.

    flink cancel -t yarn-per-job -Dyarn.application.id=<application_XXXX_YY> <jobId>

Mode application

  1. Connectez-vous au nœud maître du cluster via SSH. Pour plus d'informations, consultez Se connecter au nœud maître d'un cluster.

  2. Exécutez la commande suivante pour soumettre le job.

    flink run-application -t yarn-application /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar

    Après la soumission du job, le système renvoie un message similaire au suivant.

    [root@master-1-1(172.17.xxx.xxx) ~]# flink run-application -t yarn-application /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar
    SLF4J: Class path contains multiple SLF4J bindings.
    SLF4J: Found binding in [jar:file:/opt/apps/FLINK/flink-1.17.2-1.0.10/lib/log4j-slf4j-impl-2.17.1.jar!/org/slf4j/impl/StaticLoggerBinder.class]
    SLF4J: Found binding in [jar:file:/opt/apps/HADOOP-COMMON/hadoop-3.2.1-1.3.2-alinux3/share/hadoop/common/lib/slf4j-log4j12-1.7.jar!/org/slf4j/impl/StaticLoggerBinder.class]
    SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation.
    SLF4J: Actual binding is of type [org.apache.logging.slf4j.Log4jLoggerFactory]
    2025-06-17 10:57:05,106 INFO  org.apache.flink.yarn.cli.FlinkYarnSessionCli                [] - Found Yarn properties file under /tmp/.yarn-properties-root.
    2025-06-17 10:57:05,106 INFO  org.apache.flink.yarn.cli.FlinkYarnSessionCli                [] - Found Yarn properties file under /tmp/.yarn-properties-root.
    2025-06-17 10:57:05,233 WARN  org.apache.flink.yarn.configuration.YarnLogConfigUtil        [] - The configuration directory ('/etc/taihao-apps/flink-conf') already contains a LOG4J config file.If you want to use logback, then please delete or rename the log configuration file.
    2025-06-17 10:57:05,453 INFO  org.apache.hadoop.yarn.client.RMProxy                       [] - Connecting to ResourceManager at master-1-1.c-b9693c1xxx.cn-hangzhou.emr.aliyuncs.com/172.17.xxx.xxx:8032
    2025-06-17 10:57:05,604 INFO  org.apache.hadoop.yarn.client.AHSProxy                      [] - Connecting to Application History server at master-1-1.c-b9693xxx
    3c131faf601f.cn-hangzhou.emr.aliyuncs.com/172.17.108.111:10200
    2025-06-17 10:57:05,612 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - No path for the flink jar passed. Using the location of class org.apache.flink.yarn.YarnClusterDescriptor to locate the jar
    2025-06-17 10:57:05,724 INFO  org.apache.hadoop.conf.Configuration                        [] - found resource resource-types.xml at file:/etc/taihao-apps/hadoop-conf/resource-types.xml
    2025-06-17 10:57:05,776 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - The configured JobManager memory is 1600 MB. YARN will allocate 1664 MB to make up an integer multiple of its minimum allocation memory (128 MB, configured via 'yarn.scheduler.minimum-allocation-mb'). The extra 64 MB may not be used by Flink.
    2025-06-17 10:57:05,776 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - The configured TaskManager memory is 1728 MB. YARN will allocate 1792 MB to make up an integer multiple of its minimum allocation memory (128 MB, configured via 'yarn.scheduler.minimum-allocation-mb'). The extra 64 MB may not be used by Flink.
    2025-06-17 10:57:05,776 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Cluster specification: ClusterSpecification{masterMemoryMB=1600, taskManagerMemoryMB=1728, slotsPerTaskManager=1}
    2025-06-17 10:57:10,219 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Cannot use kerberos delegation token manager, no valid kerberos credentials provided.
    2025-06-17 10:57:10,227 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Submitting application master application_1750125819948_0004
    2025-06-17 10:57:10,271 INFO  org.apache.hadoop.yarn.client.api.impl.YarnClientImpl        [] - Submitted application application_1750125819948_0004
    2025-06-17 10:57:10,271 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Waiting for the cluster to be allocated
    2025-06-17 10:57:10,278 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Deploying cluster, current state ACCEPTED
    2025-06-17 10:57:17,825 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - YARN application has been deployed successfully.
    2025-06-17 10:57:17,825 INFO  org.apache.flink.yarn.YarnClusterDescriptor                  [] - Found Web Interface core-1-1.c-b9693c1xxx.cn-hangzhou.emr.aliyuncs.com:42563 of application 'application_1750125819948_0004'.

    Dans le message, application_1750125819948_0004 correspond à l'ID d'application YARN du job Flink soumis. Cet ID est désigné par <application_XXXX_YY> dans les sections suivantes.

  3. Exécutez la commande suivante pour afficher l'état du job.

    flink list -t yarn-application -Dyarn.application.id=<application_XXXX_YY>

    Un message similaire au suivant est renvoyé. Dans le message, 4db32b5339e6d64de2a1096c4762**** correspond à l'<jobId> du job.

    ------------------ Running/Restarting Jobs -------------------
    16.06.2025 18:20:55 : 4db32b5339e6d64de2a1096c4762**** : CarTopSpeedWindowingExample (RUNNING)

    Vous pouvez également consulter l'état du job sur l'interface web. Pour plus d'informations, consultez Afficher l'état du job sur l'interface web.

  4. Exécutez la commande suivante pour arrêter le job.

    flink cancel -t yarn-application -Dyarn.application.id=<application_XXXX_YY> <jobId>

Spécifier les configurations de job

Flink propose trois méthodes pour spécifier les configurations de job :

  • Spécifiez les valeurs de configuration directement dans le code de votre job. Pour plus d'informations, consultez Configuration Flink.

  • Lors de la soumission d'un job avec la commande flink run, utilisez l'option -D pour spécifier les valeurs de configuration. Par exemple, flink run-application -t yarn-application -D state.backend=rocksdb....

  • Spécifiez les valeurs de configuration dans le fichier /etc/taihao-apps/flink-conf/flink-conf.yaml.

Si vous ne spécifiez pas de configurations à l'aide de ces méthodes, Flink utilise les valeurs par défaut. Pour plus d'informations sur les paramètres de configuration, consultez le site officiel d'Apache Flink.

Vérifier l'état du job sur l'interface web

  1. Accédez à l'interface web.

    1. Connectez-vous à la console E-MapReduce.

    2. Dans le volet de navigation de gauche, sélectionnez EMR on ECS.

    3. Dans la barre de navigation supérieure, sélectionnez une région et un groupe de ressources selon vos besoins.

    4. Sur la page EMR on ECS, cliquez sur l'Cluster ID du cluster cible.

    5. Cliquez sur l'onglet Access Links and Ports.

    6. Sur la page Access Links and Ports, cliquez sur le lien situé dans la ligne de l'interface utilisateur YARN.

      Pour plus d'informations, consultez Accéder aux interfaces web des composants open source.

  2. Cliquez sur un ID d'application.

    Sur la page Hadoop YARN ResourceManager All Applications, recherchez l'application nommée Flink per-job cluster et cliquez sur son ID d'application (par exemple, application_1628232179762_0002).

  3. Cliquez sur le lien pour l'URL de suivi.

    Dans la section Présentation de l'application, le lien Tracking URL s'affiche sous la forme ApplicationMaster.

    La page du tableau de bord Apache Flink s'ouvre et affiche l'état du job.

Documentation connexe

Pour plus d'informations sur Flink sur YARN, consultez Apache Hadoop YARN.