將Kafka資料即時匯入到OSS等湖儲存中來降低儲存成本或者進行查詢分析是常見的使用情境。在EMR-3.37.1及之後的版本中,DataFlow叢集內建了JindoFS相關的依賴,使得您可以在DataFlow叢集中運行Flink作業,將Kafka資料以Exactly-Once語義流式寫入阿里雲OSS。本文通過樣本為您介紹如何在DataFlow叢集中編寫並運行Flink作業來滿足上述情境。
前提條件
- 已開通E-MapReduce服務和OSS服務。
- 已完成雲帳號的授權,詳情請參見角色授權。
操作流程
步驟一:準備環境
步驟二:準備JAR包
下載Demo代碼。
基於JindoFS,您可以在Flink作業中,如同HDFS一樣將資料以流式的方式寫入OSS中(路徑需要以oss://為首碼)。本樣本中使用了Flink的StreamingFileSink方法來示範開啟了檢查點(Checkpoint)之後,Flink如何以Exactly-Once語義寫入OSS。
下述程式碼片段示範了如何構建Kafka Source與OSS Sink,完整代碼您可以從GitHub連結中下載獲得。
重要JindoFS支援免密讀寫相同阿里雲帳號下的OSS儲存,因此作業中無需聲明相關AccessKey資訊。
public class OssDemoJob { public static void main(String[] args) throws Exception { ... // Check output oss dir Preconditions.checkArgument( params.get(OUTPUT_OSS_DIR).startsWith("oss://"), "outputOssDir should start with 'oss://'."); // Set up the streaming execution environment final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // Checkpoint is required env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE); String outputPath = params.get(OUTPUT_OSS_DIR); // Build Kafka source with new Source API based on 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(); // DataStream Source 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); // Compile and submit the job env.execute(); } }說明本範例程式碼片段給出了主要的樣本程式,您可以根據自身環境進行修改(例如,添加包名以及修改代碼中的Checkpoint間隔)後,進行編譯。關於如何構建Flink作業的JAR包,可以參見Flink官方文檔。如果無需任何修改,您可以直接使用 dataflow-oss-demo-1.0-SNAPSHOT.jar 包進行操作。
在命令列中,進入到下載的專案檔的根目錄下,執行以下命令打包檔案。
mvn clean package根據您pom.xml檔案中artifactId的資訊,專案對應目錄dataflow-demo/dataflow-oss-demo/target下會出現dataflow-oss-demo-1.0-SNAPSHOT.jar包。
步驟三:建立Kafka Topic並產生資料
通過SSH方式串連DataFlow叢集,詳情請參見登入叢集。
執行以下命令,建立測試所需的Topic。
kafka-topics.sh --create --bootstrap-server core-1-1:9092 \ --replication-factor 2 \ --partitions 3 \ --topic kafka-test-topic建立成功後,命令列會列印如下資訊。
Created topic kafka-test-topic.寫入資料至Kafka Topic。
在命令列中執行以下命令,進入Kafka Producer Console。
kafka-console-producer.sh --broker-list core-1-1:9092 --topic kafka-test-topic輸入五條測試資料。
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 Producer Console。
步驟四:運行Flink作業
通過SSH方式串連DataFlow叢集,詳情請參見登入叢集。
上傳打包好的dataflow-oss-demo-1.0-SNAPSHOT.jar至DataFlow叢集的根目錄下。
說明本文樣本中dataflow-oss-demo-1.0-SNAPSHOT.jar是上傳至root根目錄下,您也可以自訂上傳路徑。
執行以下命令,提交作業。
本樣本通過Per-Job Cluster模式提交作業,其他方式請參見基礎使用。
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目錄。
kafkaBrokers:指定Kafka叢集的broker,使用
core-1-1:9092即可。inputTopic:指定計劃讀取的Kafka Topic,使用在步驟三中建立的
kafka-test-topic。inputTopicGroup:指定計劃使用的Kafka Consumer Group,使用
my-group用於測試即可。
提交成功後,命令列會輸出類似如下日誌資訊。
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>為作業運行後返回的Application ID。例如,上述樣本中的application_1670236019397_0003。
步驟五:查看輸出的結果
作業正常運行後,您可以在OSS控制台查看輸出結果。
登入OSS管理主控台。
單擊建立的儲存空間。
在儲存空間的檔案管理頁面,可看到按時間命名的輸出目錄,進入目錄後可見產生的結果檔案(如 part-0-0)。
重要由於該作業為流式作業會持續運行,產生較多輸出檔案,應在完成驗證後,及時在命令列中通過
yarn application -kill <appId>命令終止該作業。
您也可以在DataFlow叢集中,通過命令列運行
hdfs dfs -cat oss://<YOUR_TARGET_BUCKET>/oss_kafka_test/<DATE_DIR>/part-0-0來展示實際儲存到OSS中的資料,樣本輸出如下。[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重要為了保證Exactly-Once語義,在Flink作業每完成一次Checkpoint(本樣本中Checkpoint間隔為30s),資料檔案才會落盤到OSS中。
此外,由於該作業為流式作業會持續運行,會產生較多輸出檔案,應在完成驗證後,及時在命令列中通過yarn application -kill <appId>命令終止該作業。