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 dari 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.
Memiliki 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 berada dalam satu Akun Alibaba Cloud yang sama, JindoFS membaca 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-nya. Lihat bagian "Stop the job" 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 dimulai 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 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-topicMasukkan 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,1662022777000Tekan 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 lain 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 yang valid 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.

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 aplikasinya adalahapplication_1670236019397_0003.
Langkah 5: Lihat output
Pekerjaan mengkomit file output setiap kali checkpoint (setiap 30 detik). Lihat output-nya, lalu hentikan pekerjaan.
Lihat output di konsol OSS:
Login ke Konsol OSS.
Di panel navigasi sebelah kiri, klik Buckets. Di halaman Buckets, klik nama bucket yang telah Anda buat.
Di halaman Objects, buka direktori output yang Anda tentukan dalam
--outputOssDir.
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
Hentikan pekerjaan:
Setelah memverifikasi output, hentikan pekerjaan untuk mencegah akumulasi file lebih lanjut:
yarn application -kill <appId>