Tutorial ini menjelaskan cara menjalankan pekerjaan streaming Flink pada kluster Dataflow E-MapReduce (EMR) yang membaca dari topik Kafka dan menulis ke Object Storage Service (OSS) dengan semantik tepat-sekali. Pekerjaan ini menggunakan JindoFS—tersedia di kluster Dataflow mulai EMR V3.37.1—untuk menulis ke OSS melalui path oss://, sama seperti menulis ke Hadoop Distributed File System (HDFS).
Prasyarat
Sebelum memulai, pastikan Anda telah:
-
Mengaktifkan EMR dan OSS pada Akun Alibaba Cloud Anda.
-
Mendapatkan izin yang diperlukan untuk Akun Alibaba Cloud atau Pengguna RAM Anda. Untuk detailnya, lihat Tetapkan peran.
Cara kerja
Pekerjaan Flink menggunakan StreamingFileSink dengan OnCheckpointRollingPolicy. Berdasarkan kebijakan ini, file output hanya dikomit ke OSS saat checkpoint selesai—bukan secara terus-menerus. Dengan interval checkpoint 30 detik, file output baru muncul kira-kira setiap 30 detik. Karena ini adalah pekerjaan streaming, pekerjaan tersebut berjalan terus-menerus dan mengumpulkan file output hingga Anda menghentikannya.
JindoFS menangani penulisan ke OSS. Jika bucket OSS dan kluster Dataflow dimiliki oleh Akun Alibaba Cloud yang sama, JindoFS membaca dari dan menulis ke bucket tersebut dalam mode tanpa password—tidak diperlukan kredensial dalam kode pekerjaan.
Karena file output terus bertambah setiap kali checkpoint, hentikan pekerjaan setelah Anda memverifikasi output. Lihat “Hentikan pekerjaan” pada Langkah 5.
Langkah 1: Siapkan lingkungan
Tutorial ini menggunakan EMR V3.43.1.
-
Buat kluster Dataflow dengan layanan Flink dan Kafka. Untuk detailnya, lihat Buat kluster.
-
Buat bucket OSS di Wilayah yang sama dengan kluster Dataflow. Untuk detailnya, lihat Buat bucket.
Langkah 2: Siapkan paket JAR
Kode contoh membuat sumber Kafka menggunakan API sumber FLIP-27 dan sink OSS yang didukung oleh StreamingFileSink. Path output harus diawali dengan oss://. Checkpoint diaktifkan dengan CheckpointingMode.EXACTLY_ONCE dan interval 30 detik—inilah yang menggerakkan OnCheckpointRollingPolicy serta menentukan seberapa sering file dikomit.
public class OssDemoJob {
public static void main(String[] args) throws Exception {
...
// Periksa direktori oss output
Preconditions.checkArgument(
params.get(OUTPUT_OSS_DIR).startsWith("oss://"),
"outputOssDir harus diawali dengan 'oss://'.");
// Siapkan environment eksekusi streaming
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Checkpoint wajib diaktifkan
env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE);
String outputPath = params.get(OUTPUT_OSS_DIR);
// Bangun sumber Kafka dengan API Source baru berbasis FLIP-27
KafkaSource<Event> kafkaSource =
KafkaSource.<Event>builder()
.setBootstrapServers(params.get(KAFKA_BROKERS_ARG))
.setTopics(params.get(INPUT_TOPIC_ARG))
.setStartingOffsets(OffsetsInitializer.latest())
.setGroupId(params.get(INPUT_TOPIC_GROUP_ARG))
.setDeserializer(new EventDeSerializationSchema())
.build();
// Sumber DataStream
DataStreamSource<Event> source =
env.fromSource(
kafkaSource,
WatermarkStrategy.<Event>forMonotonousTimestamps()
.withTimestampAssigner((event, ts) -> event.getEventTime()),
"Kafka Source");
StreamingFileSink<Event> sink =
StreamingFileSink.forRowFormat(
new Path(outputPath), new SimpleStringEncoder<Event>("UTF-8"))
.withRollingPolicy(OnCheckpointRollingPolicy.build())
.build();
source.addSink(sink);
// Kompilasi dan kirim pekerjaan
env.execute();
}
}
Potongan kode ini menunjukkan struktur program utama. Modifikasi sesuai kebutuhan—misalnya, tambahkan nama package atau sesuaikan interval checkpoint—sebelum membangun. Untuk petunjuk membangun paket JAR, lihat dokumentasi resmi Flink.
Unduh kode sumber lengkap dari dataflow-demo, atau gunakan langsung dataflow-oss-demo-1.0-SNAPSHOT.jar yang telah dibuat sebelumnya.
Untuk membuat JAR sendiri:
-
Buka direktori root proyek yang telah diunduh.
-
Jalankan perintah berikut:
mvn clean packageJAR hasil kompilasi
dataflow-oss-demo-1.0-SNAPSHOT.jarditempatkan didataflow-demo/dataflow-oss-demo/target/berdasarkanartifactIddalampom.xml.
Langkah 3: Buat topik Kafka dan hasilkan data
-
Login ke kluster Dataflow melalui SSH. Untuk detailnya, lihat Login ke kluster.
-
Buat topik Kafka untuk pengujian:
kafka-topics.sh --create --bootstrap-server core-1-1:9092 \ --replication-factor 2 \ --partitions 3 \ --topic kafka-test-topicJika berhasil, CLI akan mengembalikan:
Created topic kafka-test-topic. -
Tulis data uji ke topik tersebut.
-
Buka konsol produsen Kafka:
kafka-console-producer.sh --broker-list core-1-1:9092 --topic kafka-test-topic -
Masukkan lima catatan berikut:
1,Ken,0,1,1662022777000 1,Ken,0,2,1662022777000 1,Ken,0,3,1662022777000 1,Ken,0,4,1662022777000 1,Ken,0,5,1662022777000 -
Tekan Ctrl+C untuk keluar dari konsol produsen Kafka.
-
Langkah 4: Jalankan pekerjaan Flink
-
Login ke kluster Dataflow melalui SSH. Untuk detailnya, lihat Login ke kluster.
-
Unggah
dataflow-oss-demo-1.0-SNAPSHOT.jarke direktori root kluster Dataflow.CatatanContoh ini mengunggah ke direktori root. Pilih direktori berbeda jika diperlukan.
-
Kirim pekerjaan Flink dalam mode Per-Job:
flink run -t yarn-per-job -d -c com.alibaba.ververica.dataflow.demo.oss.OssDemoJob \ /dataflow-oss-demo-1.0-SNAPSHOT.jar \ --outputOssDir oss://xung****-flink-dlf-test/oss_kafka_test \ --kafkaBrokers core-1-1:9092 \ --inputTopic kafka-test-topic \ --inputTopicGroup my-groupParameter Deskripsi Nilai yang valid Nilai dalam contoh ini --outputOssDirPath OSS untuk menulis data. Harus diawali dengan oss://.Path OSS apa pun yang diawali dengan oss://oss://xung****-flink-dlf-test/oss_kafka_test--kafkaBrokersAlamat broker Kafka host:portcore-1-1:9092--inputTopicTopik Kafka yang akan dibaca Nama topik yang sudah ada kafka-test-topic--inputTopicGroupKelompok konsumen Kafka Nama kelompok konsumen yang valid my-groupUntuk mode pengiriman lainnya, lihat Penggunaan dasar.
Output log berikut dikembalikan:
2022-12-06 10:40:10,080 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - Submitting application master application_1670236019397_0003 2022-12-06 10:40:10,310 INFO org.apache.hadoop.yarn.client.api.impl.YarnClientImpl [] - Submitted application application_1670236019397_0003 2022-12-06 10:40:10,311 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - Waiting for the cluster to be allocated 2022-12-06 10:40:10,312 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - Deploying cluster, current state ACCEPTED 2022-12-06 10:40:16,334 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - YARN application has been deployed successfully. 2022-12-06 10:40:16,335 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - 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_1670236019397_0003 If this should not be possible, then you can also kill Flink via YARN's web interface or via: $ yarn application -kill application_1670236019397_0003 Note that killing Flink might not clean up all job artifacts and temporary files. 2022-12-06 10:40:16,335 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - Found Web Interface core-1-2xxx.cn-hangzhou.emr.aliyuncs.com:38187 of application 'application_1670236019397_0003'. Job has been submitted with JobID 6a4dbxxx488f -
Verifikasi bahwa pekerjaan sedang berjalan:
flink list -t yarn-per-job -Dyarn.application.id=<appId>Ganti
<appId>dengan ID aplikasi yang dikembalikan setelah pengiriman pekerjaan. Dalam contoh ini, ID aplikasi adalahapplication_1670236019397_0003.
Langkah 5: Lihat output
Pekerjaan mengkomit file output setiap kali checkpoint (setiap 30 detik). Lihat output, lalu hentikan pekerjaan.
Lihat output di konsol OSS:
-
Login ke Konsol OSS.
-
Di panel navigasi kiri, klik Buckets. Di halaman Buckets, klik nama bucket yang telah Anda buat.
-
Di halaman Objects, navigasi ke direktori output yang Anda tentukan di
--outputOssDir.Di halaman Objects, Anda dapat melihat direktori output yang dinamai berdasarkan timestamp. Masuk ke direktori tersebut untuk menemukan file hasil yang dihasilkan (misalnya,
part-0-0).
Lihat output dari CLI:
Jalankan perintah berikut di kluster Dataflow:
hdfs dfs -cat oss://<YOUR_TARGET_BUCKET>/oss_kafka_test/<DATE_DIR>/part-0-0
Output contoh berikut dikembalikan:
[root@master-1-xxx ~]# hdfs dfs -cat oss://xxx2022-12-06--14/part-0-0
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Hentikan pekerjaan:
Setelah memverifikasi output, hentikan pekerjaan untuk mencegah akumulasi file lebih lanjut:
yarn application -kill <appId>