Pelajari cara mengirim dan melihat pekerjaan Flink di E-MapReduce.
Latar Belakang
Layanan Flink dalam kluster Dataflow diterapkan dalam mode YARN. Anda dapat login ke kluster Dataflow melalui SSH untuk mengirim pekerjaan Flink dari command line.
Kluster Dataflow yang diterapkan dalam mode YARN mendukung pengiriman pekerjaan Flink dalam session mode, per-job cluster mode, dan application mode.
|
Mode |
Deskripsi |
Kelebihan dan Kekurangan |
|
Proses umum |
Gambar berikut menunjukkan proses umum pengiriman dan peninjauan pekerjaan Flink. Sebagai contoh, jika terjadi exception pada suatu pekerjaan yang menyebabkan TaskManager dimatikan, semua pekerjaan lain yang berjalan di TaskManager tersebut akan gagal. Selain itu, karena kluster hanya memiliki satu JobManager, beban pada JobManager meningkat seiring bertambahnya jumlah pekerjaan. |
Berdasarkan karakteristik ini, pola ini cocok untuk menerapkan pekerjaan dengan waktu startup singkat dan durasi waktu proses yang relatif pendek. |
|
Per-Job Cluster mode |
Saat menggunakan Per-Job Cluster mode, setiap kali pekerjaan Flink dikirim ke YARN, YARN akan memulai kluster Flink baru, lalu menjalankan pekerjaan tersebut. Ketika pekerjaan selesai atau dibatalkan, kluster Flink tersebut juga akan dilepas. |
Berdasarkan karakteristik di atas, mode ini biasanya cocok untuk pekerjaan berdurasi panjang. |
|
Application mode |
Saat menggunakan Application mode, setiap kali Anda mengirim Flink Application (Flink Application berisi satu atau beberapa pekerjaan) ke YARN, YARN akan memulai kluster Flink baru. Ketika Application selesai dijalankan atau dibatalkan, kluster Flink milik Application tersebut juga akan dilepas. Mode ini berbeda dari mode Per-Job dalam hal bahwa metode Jika file JAR yang dikirim berisi beberapa pekerjaan, maka semua pekerjaan tersebut akan berjalan dalam kluster milik Application tersebut. |
|
Prasyarat
Kluster Dataflow telah dibuat dalam mode Flink. Untuk informasi selengkapnya, lihat Buat kluster.
Kirim dan lihat pekerjaan Flink
Topik ini menggunakan contoh Flink TopSpeedWindowing. Contoh ini merupakan pekerjaan streaming berdurasi panjang.
Anda dapat memilih salah satu dari tiga mode berikut untuk mengirim dan melihat pekerjaan:
Session mode
-
Sambungkan ke node master kluster melalui SSH. Untuk informasi selengkapnya, lihat Login ke node master kluster.
-
Jalankan perintah berikut untuk memulai sesi YARN.
yarn-session.sh --detachedSetelah perintah berhasil dijalankan, sistem akan mengembalikan ID aplikasi, misalnya
application_1750137174986_0001. ID ini selanjutnya disebut sebagai<application_XXXX_YY>pada bagian berikutnya.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. -
Jalankan perintah berikut untuk mengirim pekerjaan.
flink run --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarSetelah pekerjaan dikirim, sistem akan mengembalikan pesan seperti berikut.
[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 3785db18d371326758d7843dd2a1xxxDalam pesan tersebut,
3785db18d371326758d7843dd2a1****adalah ID pekerjaan. ID ini selanjutnya disebut sebagai<jobId>pada bagian berikutnya. -
Jalankan perintah berikut untuk melihat status pekerjaan.
flink list -t yarn-session -Dyarn.application.id=<application_XXXX_YY>Pesan yang dikembalikan mirip dengan berikut.
------------------ Running/Restarting Jobs ------------------- 16.06.2025 18:20:55 : 3785db18d371326758d7843dd2a1**** : CarTopSpeedWindowingExample (RUNNING)Anda juga dapat melihat status pekerjaan melalui web UI. Untuk informasi selengkapnya, lihat Lihat status pekerjaan di web UI.
-
Jalankan perintah berikut untuk menghentikan pekerjaan.
flink cancel -t yarn-session -Dyarn.application.id=<application_XXXX_YY> <jobId>
Per-job cluster mode
-
Sambungkan ke node master kluster melalui SSH. Untuk informasi selengkapnya, lihat Login ke node master kluster.
-
Jalankan perintah berikut untuk mengirim pekerjaan.
flink run -t yarn-per-job --detached /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarSetelah pekerjaan dikirim, sistem akan mengembalikan pesan seperti berikut.
$ 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/rootDalam pesan tersebut,
application_1750125819948_****adalah ID aplikasi, yang selanjutnya disebut sebagai<application_XXXX_YY>pada bagian berikutnya.f5f980ac631192b02548235f1bbe****adalah ID pekerjaan, yang selanjutnya disebut sebagai<jobId>pada bagian berikutnya. -
Jalankan perintah berikut untuk melihat status pekerjaan.
flink list -t yarn-per-job -Dyarn.application.id=<application_XXXX_YY>Anda juga dapat melihat status pekerjaan melalui web UI. Untuk informasi selengkapnya, lihat Lihat status pekerjaan di web UI.
-
Jalankan perintah berikut untuk menghentikan pekerjaan.
flink cancel -t yarn-per-job -Dyarn.application.id=<application_XXXX_YY> <jobId>
Application mode
-
Sambungkan ke node master kluster melalui SSH. Untuk informasi selengkapnya, lihat Login ke node master kluster.
-
Jalankan perintah berikut untuk mengirim pekerjaan.
flink run-application -t yarn-application /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jarSetelah pekerjaan dikirim, sistem akan mengembalikan pesan seperti berikut.
[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'.Dalam pesan tersebut,
application_1750125819948_0004adalah ID aplikasi YARN dari pekerjaan Flink yang dikirim. ID ini selanjutnya disebut sebagai<application_XXXX_YY>pada bagian berikutnya. -
Jalankan perintah berikut untuk melihat status pekerjaan.
flink list -t yarn-application -Dyarn.application.id=<application_XXXX_YY>Pesan yang dikembalikan mirip dengan berikut. Dalam pesan tersebut,
4db32b5339e6d64de2a1096c4762****adalah<jobId>dari pekerjaan tersebut.------------------ Running/Restarting Jobs ------------------- 16.06.2025 18:20:55 : 4db32b5339e6d64de2a1096c4762**** : CarTopSpeedWindowingExample (RUNNING)Anda juga dapat melihat status pekerjaan melalui web UI. Untuk informasi selengkapnya, lihat Lihat status pekerjaan di web UI.
-
Jalankan perintah berikut untuk menghentikan pekerjaan.
flink cancel -t yarn-application -Dyarn.application.id=<application_XXXX_YY> <jobId>
Tentukan konfigurasi pekerjaan
Flink menyediakan tiga cara untuk menentukan konfigurasi pekerjaan:
-
Tentukan nilai konfigurasi dalam kode pekerjaan Anda. Untuk informasi selengkapnya, lihat Flink Configuration.
-
Saat mengirim pekerjaan dengan perintah
flink run, gunakan flag -D untuk menentukan nilai konfigurasi. Contohnya,flink run-application -t yarn-application -D state.backend=rocksdb.... -
Tentukan nilai konfigurasi dalam file
/etc/taihao-apps/flink-conf/flink-conf.yaml.
Jika Anda tidak menentukan konfigurasi dengan metode-metode tersebut, Flink akan menggunakan nilai default. Untuk informasi selengkapnya tentang parameter konfigurasi, lihat situs resmi Apache Flink.
Periksa status pekerjaan di web UI
-
Akses web UI.
-
Login ke Konsol E-MapReduce.
-
Pada panel navigasi kiri, pilih EMR on ECS.
-
Pada bilah navigasi atas, pilih wilayah dan kelompok sumber daya sesuai kebutuhan.
-
Pada halaman EMR on ECS, klik Cluster ID kluster target.
-
Klik tab Access Links and Ports.
-
Pada halaman Access Links and Ports, klik tautan pada baris YARN UI.
Untuk informasi selengkapnya, lihat Akses web UI komponen open-source.
-
-
Klik ID aplikasi.
Pada halaman All Applications Hadoop YARN ResourceManager, temukan aplikasi bernama Flink per-job cluster dan klik Application ID-nya (misalnya,
application_1628232179762_0002). -
Klik tautan Tracking URL.
Pada bagian Application Overview, tautan Tracking URL ditampilkan sebagai ApplicationMaster.
Halaman Apache Flink Dashboard terbuka dan menampilkan status pekerjaan.
Halaman ikhtisar Apache Flink Dashboard menampilkan informasi tentang pekerjaan yang sedang berjalan, termasuk nama pekerjaan (misalnya, CarTopSpeedWindowingExample), durasi, status tugas (RUNNING), dan jumlah Task Slot yang tersedia.
Dokumen terkait
Untuk informasi selengkapnya tentang Flink di YARN, lihat Apache Hadoop YARN.