すべてのプロダクト
Search
ドキュメントセンター

E-MapReduce:Flink を使用した Kafka データの Alibaba Cloud OSS へのストリーミング

最終更新日:Aug 22, 2026

このチュートリアルでは、Kafka トピックから読み取り、exactly-once セマンティクスで Object Storage Service (OSS) に書き込む Flink のストリーミングジョブを、E-MapReduce (EMR) Dataflow クラスター上で実行する方法を説明します。このジョブでは、EMR V3.37.1 以降の Dataflow クラスターで利用できる JindoFS を使用して、Hadoop 分散ファイルシステム (HDFS) に書き込む場合と同様に、oss:// パスを使用して OSS に書き込みます。

前提条件

開始する前に、次の項目を確認してください:

  • Alibaba Cloud アカウントで EMR と OSS が有効化されていること

  • Alibaba Cloud アカウントまたは RAM ユーザーに必要な権限が付与されていること。詳細については、「Assign roles」をご参照ください。

仕組み

Flink ジョブは、StreamingFileSinkOnCheckpointRollingPolicy を使用します。このポリシーでは、出力ファイルは継続的にコミットされるのではなく、チェックポイントが完了したときにのみ OSS にコミットされます。チェックポイント間隔が 30 秒の場合、新しい出力ファイルはおおむね 30 秒ごとに作成されます。これはストリーミングジョブであるため、停止するまで継続的に実行され、出力ファイルが蓄積されます。

OSS への書き込みは JindoFS が処理します。OSS バケットと Dataflow クラスターが同一の Alibaba Cloud アカウントに属している場合、JindoFS はパスワード不要モードでバケットの読み取りと書き込みを行うため、ジョブコード内で資格情報を指定する必要はありません。

重要

出力ファイルはチェックポイントごとに蓄積されるため、出力を確認したらジョブを停止してください。手順 5 の「ジョブの停止」を参照してください。

手順 1:環境の準備

説明

このチュートリアルでは EMR V3.43.1 を使用します。

  1. Flink サービスと Kafka サービスを含む Dataflow クラスターを作成します。詳細については、「Create a cluster」をご参照ください。

  2. Dataflow クラスターと同じリージョンに OSS バケットを作成します。詳細については、「Create buckets」をご参照ください。

手順 2:JAR パッケージの準備

サンプルコードは、FLIP-27 ソース API を使用して Kafka ソースと、StreamingFileSink で実装された OSS シンクを作成します。出力パスは oss:// で始まる必要があります。チェックポイント機能は CheckpointingMode.EXACTLY_ONCE と 30 秒の間隔で有効になります。これが OnCheckpointRollingPolicy を駆動し、ファイルがコミットされる頻度を決定します。

public class OssDemoJob {

    public static void main(String[] args) throws Exception {
        ...

        // 出力 OSS ディレクトリの確認
        Preconditions.checkArgument(
                params.get(OUTPUT_OSS_DIR).startsWith("oss://"),
                "outputOssDir should start with 'oss://'.");

        // ストリーミング実行環境のセットアップ
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // チェックポイントは必須
        env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE);

        String outputPath = params.get(OUTPUT_OSS_DIR);

        // FLIP-27 に基づく新しい Source API で Kafka ソースを構築
        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();
        // 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);

        // ジョブをコンパイルしてサブミット
        env.execute();
    }
}
説明

このスニペットはメインプログラムの構成を示しています。ビルド前に必要に応じて変更してください。例えば、パッケージ名を追加したり、チェックポイント間隔を調整したりできます。JAR パッケージのビルド手順については、「Flink official documentation」をご参照ください。

完全なソースは dataflow-demo からダウンロードできます。または、事前ビルド済みの dataflow-oss-demo-1.0-SNAPSHOT.jar を直接使用できます。

JAR を自分でビルドする場合:

  1. ダウンロードしたプロジェクトのルートディレクトリに移動します。

  2. 次のコマンドを実行します:

    mvn clean package

    出力 JAR の dataflow-oss-demo-1.0-SNAPSHOT.jar は、pom.xmlartifactId に基づき、dataflow-demo/dataflow-oss-demo/target/ に配置されます。

手順 3:Kafka トピックの作成とデータの生成

  1. SSH 経由で Dataflow クラスターにログインします。詳細については、「Log on to a cluster」をご参照ください。

  2. テスト用の Kafka トピックを作成します:

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

    成功した場合、CLI に次のメッセージが表示されます:

    Created topic kafka-test-topic.
  3. テストデータをトピックに書き込みます。

    1. Kafka プロデューサーコンソールを開きます:

      kafka-console-producer.sh --broker-list core-1-1:9092 --topic  kafka-test-topic
    2. 次の 5 件のレコードを入力します:

      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. Ctrl+C を押して Kafka プロデューサーコンソールを終了します。

手順 4:Flink ジョブの実行

  1. SSH 経由で Dataflow クラスターにログインします。詳細については、「Log on to a cluster」をご参照ください。

  2. dataflow-oss-demo-1.0-SNAPSHOT.jar を Dataflow クラスターのルートディレクトリにアップロードします。

    説明

    この例ではルートディレクトリにアップロードします。必要に応じて別のディレクトリを選択してください。

  3. Per-Job モードで Flink ジョブをサブミットします:

    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
    パラメーター 説明 有効な値 この例の値
    --outputOssDir データの書き込み先となる OSS パス。oss:// で始まる必要があります。 oss:// で始まる任意の有効な OSS パス oss://xung****-flink-dlf-test/oss_kafka_test
    --kafkaBrokers Kafka ブローカーのアドレス host:port core-1-1:9092
    --inputTopic 読み取り元の Kafka トピック 既存の任意のトピック名 kafka-test-topic
    --inputTopicGroup Kafka コンシューマーグループ 有効な任意のコンシューマーグループ名 my-group

    その他のサブミットモードについては、「Basic usage」をご参照ください。

    次のログ出力が表示されます:

  4. ジョブが実行中であることを確認します:

    flink list -t yarn-per-job -Dyarn.application.id=<appId>

    <appId> を、ジョブのサブミット後に返されたアプリケーション ID に置き換えてください。この例では、アプリケーション ID は application_1670236019397_0003 です。

手順 5:出力の確認

ジョブはチェックポイントごと (30 秒ごと) に出力ファイルをコミットします。出力を確認してから、ジョブを停止してください。

OSS コンソールで出力を確認:

  1. OSS console にログインします。

  2. 左側のナビゲーションウィンドウで [Buckets] をクリックします。[Buckets] ページで、作成したバケットの名前をクリックします。

  3. [Objects] ページで、--outputOssDir で指定した出力ディレクトリに移動します。

    [Objects] ページでは、タイムスタンプ名の出力ディレクトリを確認できます。ディレクトリ内に移動して、生成された結果ファイル (例:part-0-0) を確認します。

CLI で出力を確認:

Dataflow クラスター上で次のコマンドを実行します:

hdfs dfs -cat oss://<YOUR_TARGET_BUCKET>/oss_kafka_test/<DATE_DIR>/part-0-0

次のサンプル出力が表示されます:

[root@master-1-xxx ~]# hdfs dfs -cat oss://xung****-flink-dlf-test/oss_kafka_test/2022-12-06--10-40/part-0-0
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,2,1662022777000
Sink to OSS: 1,Ken,0,3,1662022777000
Sink to OSS: 1,Ken,0,4,1662022777000
Sink to OSS: 1,Ken,0,5,1662022777000

ジョブの停止:

出力を確認したら、ファイルがさらに蓄積されないようにジョブを停止します:

yarn application -kill <appId>