このチュートリアルでは、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 ジョブは、StreamingFileSink と OnCheckpointRollingPolicy を使用します。このポリシーでは、出力ファイルは継続的にコミットされるのではなく、チェックポイントが完了したときにのみ OSS にコミットされます。チェックポイント間隔が 30 秒の場合、新しい出力ファイルはおおむね 30 秒ごとに作成されます。これはストリーミングジョブであるため、停止するまで継続的に実行され、出力ファイルが蓄積されます。
OSS への書き込みは JindoFS が処理します。OSS バケットと Dataflow クラスターが同一の Alibaba Cloud アカウントに属している場合、JindoFS はパスワード不要モードでバケットの読み取りと書き込みを行うため、ジョブコード内で資格情報を指定する必要はありません。
出力ファイルはチェックポイントごとに蓄積されるため、出力を確認したらジョブを停止してください。手順 5 の「ジョブの停止」を参照してください。
手順 1:環境の準備
このチュートリアルでは EMR V3.43.1 を使用します。
-
Flink サービスと Kafka サービスを含む Dataflow クラスターを作成します。詳細については、「Create a cluster」をご参照ください。
-
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 を自分でビルドする場合:
-
ダウンロードしたプロジェクトのルートディレクトリに移動します。
-
次のコマンドを実行します:
mvn clean package出力 JAR の
dataflow-oss-demo-1.0-SNAPSHOT.jarは、pom.xmlのartifactIdに基づき、dataflow-demo/dataflow-oss-demo/target/に配置されます。
手順 3:Kafka トピックの作成とデータの生成
-
SSH 経由で Dataflow クラスターにログインします。詳細については、「Log on to a cluster」をご参照ください。
-
テスト用の 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. -
テストデータをトピックに書き込みます。
-
Kafka プロデューサーコンソールを開きます:
kafka-console-producer.sh --broker-list core-1-1:9092 --topic kafka-test-topic -
次の 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 -
Ctrl+C を押して Kafka プロデューサーコンソールを終了します。
-
手順 4:Flink ジョブの実行
-
SSH 経由で Dataflow クラスターにログインします。詳細については、「Log on to a cluster」をご参照ください。
-
dataflow-oss-demo-1.0-SNAPSHOT.jarを Dataflow クラスターのルートディレクトリにアップロードします。説明この例ではルートディレクトリにアップロードします。必要に応じて別のディレクトリを選択してください。
-
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--kafkaBrokersKafka ブローカーのアドレス host:portcore-1-1:9092--inputTopic読み取り元の Kafka トピック 既存の任意のトピック名 kafka-test-topic--inputTopicGroupKafka コンシューマーグループ 有効な任意のコンシューマーグループ名 my-groupその他のサブミットモードについては、「Basic usage」をご参照ください。
次のログ出力が表示されます:
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 -
ジョブが実行中であることを確認します:
flink list -t yarn-per-job -Dyarn.application.id=<appId><appId>を、ジョブのサブミット後に返されたアプリケーション ID に置き換えてください。この例では、アプリケーション ID はapplication_1670236019397_0003です。
手順 5:出力の確認
ジョブはチェックポイントごと (30 秒ごと) に出力ファイルをコミットします。出力を確認してから、ジョブを停止してください。
OSS コンソールで出力を確認:
-
OSS console にログインします。
-
左側のナビゲーションウィンドウで [Buckets] をクリックします。[Buckets] ページで、作成したバケットの名前をクリックします。
-
[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>