All Products
Search
Document Center

E-MapReduce:Alirkan data Kafka ke Alibaba Cloud OSS menggunakan Flink

Last Updated:Aug 22, 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 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.

Penting

Karena file output terus bertambah setiap kali checkpoint, hentikan pekerjaan setelah Anda memverifikasi output. Lihat “Hentikan pekerjaan” 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 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();
    }
}
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 sumber 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 berbeda 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
    Parameter Deskripsi Nilai yang valid Nilai dalam contoh ini
    --outputOssDir Path OSS untuk menulis data. Harus diawali dengan oss://. Path OSS apa pun yang diawali dengan oss:// oss://xung****-flink-dlf-test/oss_kafka_test
    --kafkaBrokers Alamat broker Kafka host:port core-1-1:9092
    --inputTopic Topik Kafka yang akan dibaca Nama topik yang sudah ada kafka-test-topic
    --inputTopicGroup Kelompok konsumen Kafka Nama kelompok konsumen yang valid my-group

    Untuk mode pengiriman lainnya, lihat Penggunaan dasar.

    Output log berikut dikembalikan:

  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 aplikasi adalah application_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:

  1. Login ke Konsol OSS.

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

  3. 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>