Flink DataStream API は、複雑なビジネスロジックとデータ処理を扱うカスタムのデータ変換、操作、およびオペレーターを作成するための柔軟なプログラミングモデルを提供します。
Apache Flink との互換性
Realtime Compute for Apache Flink がサポートする DataStream API は、オープンソースの Apache Flink バージョンと完全な互換性があります。詳細については、「Apache Flink とは」および「Flink DataStream API プログラミングガイド」をご参照ください。
前提条件
-
IntelliJ IDEA などの統合開発環境 (IDE) がインストールされていること。
-
Maven 3.6.3 以降がインストールされていること。
-
ジョブ開発は JDK 8 および JDK 11 のみをサポートします。
-
Realtime Compute for Apache Flink コンソールでデプロイして実行する前に、ローカルで JAR ジョブを開発してください。
事前準備
この例ではデータソースコネクタを使用するため、事前に必要なデータソースを準備してください。
-
この例では、データソースとして ApsaraMQ for Kafka (2.6.2) および ApsaraDB RDS for MySQL (8.0) を使用します。
-
パブリックネットワークアクセスまたは VPC 間のアクセスを必要とする自己管理データソースがある場合は、「ネットワーク接続オプション」をご参照ください。
-
ApsaraMQ for Kafka データソースがない場合は、インスタンスを購入してデプロイしてください。詳細については、「ステップ2:インスタンスの購入とデプロイ」をご参照ください。インスタンスをデプロイする際は、Realtime Compute for Apache Flink ワークスペースと同じ VPC 内にあることを確認してください。
-
ApsaraDB RDS for MySQL データソースがない場合は、ApsaraDB RDS for MySQL インスタンスを購入してください。詳細については、「ステップ1:ApsaraDB RDS for MySQL インスタンスの作成とデータベースの設定」をご参照ください。インスタンスを購入する際は、Realtime Compute for Apache Flink ワークスペースと同じリージョンおよび VPC 内にあることを確認してください。
ジョブの開発
Flink 環境の依存関係の設定
JAR の依存関係の競合を避けるために、以下のガイドラインに従ってください。
-
${flink.version}は、ジョブ実行用の Flink バージョンを指定します。このバージョンは、ジョブデプロイページで選択した VVR エンジンの Flink バージョンと一致している必要があります。たとえば、ジョブデプロイページでvvr-8.0.4-flink-1.17エンジンを選択した場合、対応する Flink バージョンは1.17.1です。VVR エンジンのバージョンに関する詳細については、「現在のジョブの Flink バージョンを確認する方法」をご参照ください。 -
Flink の依存関係については、
<scope>provided</scope>を追加して、スコープをprovidedに設定します。これには主に、org.apache.flinkグループ内のflink-で始まるコネクタ以外の依存関係が含まれます。 -
Flink のソースコードでは、@Public または @PublicEvolving で明示的にアノテーションが付けられたメソッドのみがパブリック API です。Realtime Compute for Apache Flink は、これらのメソッドの互換性のみを保証します。
-
組み込みの Flink コネクタの DataStream API を使用する場合は、provided 依存関係を使用します。
以下は基本的な Flink の依存関係です。ロギングの依存関係も追加する必要がある場合があります。依存関係の完全なリストについては、このトピックの最後にある「完全なコード例」をご参照ください。
Flink の依存関係
コネクタの依存関係と使用方法
DataStream API でデータを読み書きするには、DataStream コネクタを使用して Realtime Compute for Apache Flink に接続します。Maven 中央リポジトリは、開発中に使用できる VVR DataStream コネクタを提供しています。
「サポートされているコネクタ」で DataStream API をサポートすると指定されているコネクタのみを使用してください。DataStream API のサポートが明記されていないコネクタは、将来的にインターフェイスやパラメータが変更される可能性があるため、使用しないでください。
コネクタは、以下のいずれかの方法で使用できます。
(推奨) 追加の依存関係としてアップロード
-
ジョブの pom.xml ファイルに、必要なコネクタを "provided" スコープのプロジェクト依存関係として追加します。完全な依存関係ファイルについては、このトピックの最後にある「完全なコード例」をご参照ください。
説明-
${vvr.version}はコネクタのバージョンです。ジョブのランタイムエンジンバージョンと互換性がある必要があります。たとえば、ジョブがvvr-8.0.4-flink-1.17エンジンで実行される場合、対応する Flink バージョンは1.17.1です。最新のエンジンの使用を推奨します。特定のバージョンに関する情報については、「エンジン」をご参照ください。 -
コネクタの JAR パッケージは追加の依存関係として追加されるため、アプリケーション JAR にバンドルする必要はありません。したがって、スコープを
providedに設定する必要があります。
<!-- Kafkaコネクタの依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency> <!-- MySQLコネクタの依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency> -
-
新しいコネクタを開発したり、既存のコネクタの機能を拡張したりする必要がある場合、プロジェクトは共通のコネクタパッケージ
flink-connector-baseまたはververica-connector-commonにも依存する必要があります。<!-- Flinkコネクタのパブリックインターフェイスの基本依存関係 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Alibaba Cloudコネクタのパブリックインターフェイスの基本依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> -
DataStream の接続設定とコード例については、対応する DataStream コネクタのドキュメントをご参照ください。
DataStream API をサポートするコネクタのリストについては、「サポートされているコネクタ」をご参照ください。
-
ジョブをデプロイし、[追加の依存関係] フィールドに対応するコネクタ JAR パッケージを追加します。詳細については、「JAR ジョブのデプロイ」をご参照ください。カスタムコネクタまたは Realtime Compute for Apache Flink が提供するコネクタをアップロードできます。コネクタをダウンロードするには、「コネクタ」をご参照ください。
たとえば、
ververica-connector-mysql-1.17-vvr-8.0.4-1.jarとververica-connector-kafka-1.17-vvr-8.0.4-1.jarコネクタ JAR パッケージ、およびconfig.properties設定ファイルをアップロードします。
ジョブ JAR でのパッケージ化
-
ジョブの pom.xml ファイルに、必要なコネクタをプロジェクト依存関係として追加します。次のコードは、Kafka および MySQL コネクタの例を示しています。
説明-
${vvr.version}はコネクタのバージョンです。ジョブのランタイム環境のエンジンバージョンと互換性がある必要があります。たとえば、ジョブがvvr-8.0.4-flink-1.17エンジンバージョンで実行される場合、対応する Flink バージョンは1.17.1です。最新のエンジンの使用を推奨します。詳細については、「エンジン」をご参照ください。 -
コネクタをプロジェクト依存関係としてジョブ JAR にパッケージ化する場合、それらはデフォルトスコープ (compile) である必要があります。
<!-- Kafkaコネクタの依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> </dependency> <!-- MySQLコネクタの依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> </dependency> -
-
新しいコネクタを開発したり、既存のコネクタの機能を拡張したりする必要がある場合、プロジェクトには共通のコネクタパッケージ
flink-connector-baseまたはververica-connector-commonも必要です。<!-- Flinkコネクタのパブリックインターフェイスの基本依存関係 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Alibaba Cloudコネクタのパブリックインターフェイスの基本依存関係 --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> -
DataStream の接続設定とコード例については、対応する DataStream コネクタのドキュメントをご参照ください。
DataStream API をサポートするコネクタのリストについては、「サポートされているコネクタ」をご参照ください。
OSS からの追加依存関係の読み取り
Flink JAR ジョブは、main メソッドからローカル設定ファイルを読み取ることはできません。代わりに、設定ファイルをワークスペースの OSS バケットにアップロードし、デプロイ時に追加依存関係として追加し、実行時に読み取ります。次のセクションで例を示します。
-
コードに認証情報をハードコーディングするのを避けるため、config.properties という名前の設定ファイルを作成します。
# Kafka bootstrapServers=host1:9092,host2:9092,host3:9092 inputTopic=topic groupId=groupId # MySQL database.url=jdbc:mysql://localhost:3306/my_database database.username=username database.password=password -
JAR ジョブで、OSS バケットに保存されている config.properties ファイルを読み取るコードを使用します。
方法1:ワークスペースバケットからの読み取り
-
Realtime Compute for Apache Flink コンソールの左側メニューで、[アーティファクト] ページに移動し、ファイルをアップロードします。
-
実行時に、[追加の依存関係] フィールドに追加されたファイルは、ジョブが実行される Pod の /flink/usrlib ディレクトリにロードされます。
-
次のコードは、この設定ファイルを読み取る方法の例です。
Properties properties = new Properties(); Map<String,String> configMap = new HashMap<>(); try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) { // プロパティファイルをロードします。 properties.load(input); // プロパティ値を取得します。 configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ; configMap.put("inputTopic",properties.getProperty("inputTopic")); configMap.put("groupId",properties.getProperty("groupId")); configMap.put("url",properties.getProperty("database.url")) ; configMap.put("username",properties.getProperty("database.username")); configMap.put("password",properties.getProperty("database.password")); } catch (IOException ex) { ex.printStackTrace(); }
方法2:承認済みバケットからの読み取り
-
設定ファイルをターゲットの OSS バケットにアップロードします。
-
OSSClient を使用して、OSS から直接ファイルを読み取ります。詳細については、「ストリーミングダウンロード」および「アクセス資格情報の管理」をご参照ください。次のコードは例を示しています。
OSS ossClient = new OSSClientBuilder().build("Endpoint", "AccessKeyId", "AccessKeySecret"); try (OSSObject ossObject = ossClient.getObject("examplebucket", "exampledir/config.properties"); BufferedReader reader = new BufferedReader(new InputStreamReader(ossObject.getObjectContent()))) { // ファイルを読み取り、処理します... } finally { if (ossClient != null) { ossClient.shutdown(); } }
-
ビジネスロジックの作成
-
外部データソースを Flink データストリームプログラムに統合できます。
ウォーターマークは、時間セマンティクスに基づく Flink の計算戦略であり、タイムスタンプと共に使用されることが多いため、この例ではウォーターマーク戦略を使用しません。詳細については、「ウォーターマーク戦略」をご参照ください。// 外部データソースを Flink データストリームプログラムに統合します。 // WatermarkStrategy.noWatermarks() は、ウォーターマーク戦略を使用しないことを示します。 DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source"); -
オペレーター変換では、この例では
DataStream<String>をDataStream<Student>に変換します。より複雑なオペレーター変換と処理方法については、「Flink オペレーター」をご参照ください。// データ構造を Student に変換するオペレーター。 DataStream<Student> source = stream .map(new MapFunction<String, Student>() { @Override public Student map(String s) throws Exception { // データはカンマで区切られます。 String[] data = s.split(","); return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2])); } }).filter(student -> student.score >=60); // スコアが 60 以上のデータをフィルタリングします。
ジョブのパッケージ化
maven-shade-plugin を使用してジョブをパッケージ化します。
-
コネクタを追加の依存関係として追加する場合、ジョブをパッケージ化する際にコネクタの依存関係のスコープが
providedに設定されていることを確認してください。 -
コネクタをジョブ JAR にパッケージ化する場合は、デフォルトスコープ (compile) を使用します。
ジョブのテストとデプロイ
-
デフォルトでは、Realtime Compute for Apache Flink はパブリックインターネットにアクセスできないため、ローカル環境での直接テストはできません。単体テストを個別に行うことを推奨します。詳細については、「コネクタを使用したジョブのローカルでの実行とデバッグ」をご参照ください。
-
JAR ジョブをデプロイするには、「JAR ジョブのデプロイ」をご参照ください。
説明-
デプロイ中に、コネクタを追加の依存関係として使用する場合は、関連する JAR パッケージを必ずアップロードしてください。
-
設定ファイルを読み取る必要がある場合は、それを追加の依存関係としてアップロードする必要もあります。
-
完全なコード例
この例では、ApsaraMQ for Kafka ソースからのデータを処理し、結果を ApsaraDB RDS for MySQL シンクに書き込みます。この例は参考用です。コードスタイルと品質に関するガイドラインの詳細については、「コードスタイルと品質ガイド」をご参照ください。
この例では、チェックポイント、TTL、再起動戦略などのランタイムパラメータの設定は省略されています。これらの設定は、ジョブがデプロイされた後、[デプロイ詳細] ページで設定できます。コード内の設定は優先度が高いため、将来の更新を簡素化するためにデプロイ後に設定することを推奨します。詳細については、「ジョブデプロイ情報の設定」をご参照ください。
FlinkDemo.java
package com.aliyun;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class FlinkDemo {
// データ構造を定義します。
public static class Student {
public int id;
public String name;
public int score;
public Student(int id, String name, int score) {
this.id = id;
this.name = name;
this.score = score;
}
}
public static void main(String[] args) throws Exception {
// Flink 実行環境を作成します。
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
Map<String,String> configMap = new HashMap<>();
try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
// プロパティファイルをロードします。
properties.load(input);
// プロパティ値を取得します。
configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ;
configMap.put("inputTopic",properties.getProperty("inputTopic"));
configMap.put("groupId",properties.getProperty("groupId"));
configMap.put("url",properties.getProperty("database.url")) ;
configMap.put("username",properties.getProperty("database.username"));
configMap.put("password",properties.getProperty("database.password"));
} catch (IOException ex) {
ex.printStackTrace();
}
// Kafka ソースをビルドします
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers(configMap.get("bootstrapServers"))
.setTopics(configMap.get("inputTopic"))
.setStartingOffsets(OffsetsInitializer.latest())
.setGroupId(configMap.get("groupId"))
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
.build();
// 外部データソースを Flink データストリームプログラムに統合します。
// WatermarkStrategy.noWatermarks() は、ウォーターマーク戦略を使用しないことを示します。
DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
// スコアが 60 以上のものをフィルタリングします。
DataStream<Student> source = stream
.map(new MapFunction<String, Student>() {
@Override
public Student map(String s) throws Exception {
String[] data = s.split(",");
return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
}
}).filter(Student -> Student.score >=60);
source.addSink(JdbcSink.sink("INSERT IGNORE INTO student (id, username, score) VALUES (?, ?, ?)",
new JdbcStatementBuilder<Student>() {
public void accept(PreparedStatement ps, Student data) {
try {
ps.setInt(1, data.id);
ps.setString(2, data.name);
ps.setInt(3, data.score);
} catch (SQLException e) {
throw new RuntimeException(e);
}
}
},
new JdbcExecutionOptions.Builder()
.withBatchSize(5) // バッチ書き込みごとのレコード数。
.withBatchIntervalMs(2000) // バッチ間隔 (ミリ秒)。
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl(configMap.get("url"))
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername(configMap.get("username"))
.withPassword(configMap.get("password"))
.build()
)).name("Sink MySQL");
env.execute("Flink Demo");
}
}
pom.xml
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.aliyun</groupId>
<artifactId>FlinkDemo</artifactId>
<version>1.0-SNAPSHOT</version>
<name>FlinkDemo</name>
<packaging>jar</packaging>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<flink.version>1.17.1</flink.version>
<vvr.version>1.17-vvr-8.0.4-1</vvr.version>
<target.java.version>1.8</target.java.version>
<maven.compiler.source>${target.java.version}</maven.compiler.source>
<maven.compiler.target>${target.java.version}</maven.compiler.target>
<log4j.version>2.14.1</log4j.version>
</properties>
<dependencies>
<!-- Apache Flink の依存関係 -->
<!-- これらの依存関係は JAR ファイルにパッケージ化すべきではないため、「provided」に設定されています。 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<!-- コネクタの依存関係をここに追加します。これらは 'provided' スコープに設定する必要があります。 -->
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-kafka</artifactId>
<version>${vvr.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
<scope>provided</scope>
</dependency>
<!-- 実行時にコンソール出力を生成するためのロギングフレームワークを追加します。 -->
<!-- デフォルトでは、これらの依存関係はアプリケーション JAR から除外されます。 -->
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-slf4j-impl</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-api</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
</dependencies>
<build>
<plugins>
<!-- Java コンパイラ -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.11.0</version>
<configuration>
<source>${target.java.version}</source>
<target>${target.java.version}</target>
</configuration>
</plugin>
<!-- maven-shade-plugin を使用して、必要なすべての依存関係を含む fat JAR を作成します。 -->
<!-- プログラムのエントリポイントが変更された場合は、<mainClass> の値を変更してください。 -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<!-- 不要な依存関係を除外します。 -->
<configuration>
<artifactSet>
<excludes>
<exclude>org.apache.flink:force-shading</exclude>
<exclude>com.google.code.findbugs:jsr305</exclude>
<exclude>org.slf4j:*</exclude>
<exclude>org.apache.logging.log4j:*</exclude>
</excludes>
</artifactSet>
<filters>
<filter>
<!-- META-INF ディレクトリから署名をコピーしないでください。
コピーすると、JAR ファイルの使用時にセキュリティ例外がスローされる可能性があります。 -->
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.aliyun.FlinkDemo</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
関連ドキュメント
-
DataStream API をサポートするコネクタのリストについては、「サポートされているコネクタ」をご参照ください。
-
Flink JAR ジョブ開発ワークフローの完全なエンドツーエンドの例については、「Flink JAR ジョブ」をご参照ください。
-
Realtime Compute for Apache Flink は、SQL および Python ジョブもサポートしています。開発ガイダンスについては、「ジョブ開発の概要」および「Python ジョブの開発」をご参照ください。