All Products
Search
Document Center

E-MapReduce:Gunakan Flink untuk menulis data Kafka ke OSS dalam mode streaming

Last Updated:Mar 26, 2026

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.

Penting

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

Catatan

Tutorial ini menggunakan EMR V3.43.1.

  1. Buat kluster Dataflow dengan layanan Flink dan Kafka. Untuk detailnya, lihat Buat kluster.

  2. 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();
    }
}
Catatan

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:

  1. Buka direktori root proyek yang telah diunduh.

  2. Jalankan perintah berikut:

    mvn clean package

    JAR hasil kompilasi dataflow-oss-demo-1.0-SNAPSHOT.jar ditempatkan di dataflow-demo/dataflow-oss-demo/target/ berdasarkan artifactId dalam pom.xml.

Langkah 3: Buat topik Kafka dan hasilkan data

  1. Login ke kluster Dataflow melalui SSH. Untuk detailnya, lihat Login ke kluster.

  2. Buat topik Kafka untuk pengujian:

    kafka-topics.sh --create  --bootstrap-server core-1-1:9092 \
        --replication-factor 2  \
        --partitions 3  \
        --topic kafka-test-topic

    Jika berhasil, CLI akan mengembalikan:

    Created topic kafka-test-topic.
  3. Tulis data uji ke topik tersebut.

    1. Buka konsol produsen Kafka:

      kafka-console-producer.sh --broker-list core-1-1:9092 --topic  kafka-test-topic
    2. 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
    3. Tekan Ctrl+C untuk keluar dari konsol produsen Kafka.

Langkah 4: Jalankan pekerjaan Flink

  1. Login ke kluster Dataflow melalui SSH. Untuk detailnya, lihat Login ke kluster.

  2. Unggah dataflow-oss-demo-1.0-SNAPSHOT.jar ke direktori root kluster Dataflow.

    Catatan

    Contoh ini mengunggah ke direktori root. Pilih direktori lain jika diperlukan.

  3. 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-group
    ParameterDeskripsiNilai yang validNilai 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 Kafkahost:portcore-1-1:9092
    --inputTopicTopik Kafka yang akan dibacaNama topik yang sudah adakafka-test-topic
    --inputTopicGroupKelompok konsumen KafkaNama kelompok konsumen yang validmy-group

    Untuk mode pengiriman lainnya, lihat Penggunaan dasar.

    result

  4. 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 adalah application_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:

  1. Login ke Konsol OSS.

  2. Di panel navigasi sebelah kiri, klik Buckets. Di halaman Buckets, klik nama bucket yang telah Anda buat.

  3. Di halaman Objects, buka direktori output yang Anda tentukan dalam --outputOssDir.

    OSS results

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
OSS示例

Hentikan pekerjaan:

Setelah memverifikasi output, hentikan pekerjaan untuk mencegah akumulasi file lebih lanjut:

yarn application -kill <appId>