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. |
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é. |
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 Si le fichier JAR soumis contient plusieurs jobs, tous ces jobs s'exécuteront dans le cluster de l'application. |
|
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
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
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.
-
Exécutez la commande suivante pour démarrer une session YARN.
yarn-session.sh --detachedUne 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. -
Exécutez la commande suivante pour soumettre le job.
flink run --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarAprè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 3785db18d371326758d7843dd2a1xxxDans le message,
3785db18d371326758d7843dd2a1****correspond à l'ID du job. Cet ID est désigné par<jobId>dans les sections suivantes. -
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.
-
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
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.
-
Exécutez la commande suivante pour soumettre le job.
flink run -t yarn-per-job --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarAprè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/rootDans 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. -
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.
-
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
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.
-
Exécutez la commande suivante pour soumettre le job.
flink run-application -t yarn-application /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarAprè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_0004correspond à l'ID d'application YARN du job Flink soumis. Cet ID est désigné par<application_XXXX_YY>dans les sections suivantes. -
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.
-
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
-
Accédez à l'interface web.
Connectez-vous à la console E-MapReduce.
Dans le volet de navigation de gauche, sélectionnez EMR on ECS.
Dans la barre de navigation supérieure, sélectionnez une région et un groupe de ressources selon vos besoins.
Sur la page EMR on ECS, cliquez sur l'Cluster ID du cluster cible.
Cliquez sur l'onglet Access Links and Ports.
-
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.
-
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). -
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.
La page de présentation du tableau de bord Apache Flink affiche les informations relatives au job en cours d'exécution, notamment le nom du job (par exemple, CarTopSpeedWindowingExample), la durée, l'état des tâches (RUNNING) et le nombre de slots de tâches disponibles.
Documentation connexe
Pour plus d'informations sur Flink sur YARN, consultez Apache Hadoop YARN.