このトピックでは、MySQL コネクタを使用する DataStream ジョブのデバッグと実行方法について説明します。
MySQL CDC DataStream API
DataStream API を使用してデータを読み書きするには、対応する DataStream コネクタを使用する必要があります。DataStream コネクタの設定方法の詳細については、「DataStream コネクタのセットアップ」をご参照ください。
DataStream API プログラムを作成し、MySqlSource を使用します。以下の例は、コードと pom 依存関係を示しています。
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
public class MySqlSourceExample {
public static void main(String[] args) throws Exception {
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("yourHostname")
.port(yourPort)
.databaseList("yourDatabaseName") // キャプチャ対象のデータベースを設定
.tableList("yourDatabaseName.yourTableName") // キャプチャ対象のテーブルを設定
.username("yourUsername")
.password("yourPassword")
.deserializer(new JsonDebeziumDeserializationSchema()) // SourceRecord を JSON 文字列に変換
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// チェックポイントを有効化
env.enableCheckpointing(3000);
env
.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
// 4 つの並列ソースタスクを設定
.setParallelism(4)
.print().setParallelism(1); // sink の並列度を 1 に設定してメッセージの順序を維持
env.execute("Print MySQL Snapshot + Binlog");
}
}<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-core</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-connector-base</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-common</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>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
</dependency>MySqlSource をビルドする際、コード内で以下のパラメーターを指定する必要があります。
|
パラメーター |
説明 |
|
hostname |
MySQL データベースの IP アドレスまたはホスト名。 |
|
port |
MySQL データベースサービスのポート番号。 |
|
databaseList |
MySQL データベースの名前。 説明
データベース名では、正規表現を使用して複数のデータベースからデータを読み取ることができます。すべてのデータベースを照合するには、 |
|
username |
MySQL データベースサービスのユーザー名。 |
|
password |
MySQL データベースサービスのパスワード。 |
|
deserializer |
SourceRecord タイプのレコードを指定されたタイプに変換するデシリアライザ。有効な値は次のとおりです:
|
pom の依存関係では以下のパラメーターを指定する必要があります。
|
${vvr.version} |
リアルタイムコンピューティング for Apache Flink のエンジンバージョン。例: 説明
ホットフィックスバージョンが随時リリースされ、これらの更新が他のチャネルを通じて告知されない場合があるため、Maven に表示されるバージョン番号を正しいものとして使用してください。 |
|
${flink.version} |
Apache Flink のバージョン。例: 重要
ランタイムでの非互換性の問題を回避するために、Realtime Compute for Apache Flink のエンジンバージョンに対応する Apache Flink バージョンを使用してください。バージョンマッピングの詳細については、「エンジン」をご参照ください。 |
DataStream のデバッグソリューション
DataStream API プログラムを作成し、MySqlSource を使用します。以下にコードの例を示します。
ローカルデバッグの場合、必要な JAR ファイルをダウンロードし、依存関係を設定する必要があります。詳細については、「コネクタを含むジョブをローカルで実行およびデバッグする」をご参照ください。このトピックでは、サンプルの Maven プロジェクトを提供しています。詳細については、「MysqlCDCDemo.zip」をご参照ください。
package com.alibaba.realtimecompute;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.configuration.Configuration;
public class MysqlCDCDemo {
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
conf.setString("pipeline.classpaths", "file://" + "absolute path to the MySQL uber JAR"); // 依存関係を設定
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("hostname")
.port(3306)
.databaseList("test_db") // キャプチャ対象のデータベースを設定
.tableList("test_db.test_table") // キャプチャ対象のテーブルを設定
.username("username")
.password("password")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
env.enableCheckpointing(3000);
env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
.setParallelism(4)
.print().setParallelism(1);
env.execute("Print MySQL Snapshot + Binlog");
}
}
MySqlSource をビルドする際、コード内で以下のパラメーターを指定する必要があります。
|
パラメーター |
説明 |
|
hostname |
MySQL データベースの IP アドレスまたはホスト名。 |
|
port |
MySQL データベースサービスのポート番号。 |
|
databaseList |
MySQL データベースの名前。 説明
データベース名は、複数のデータベースからデータを読み取るための正規表現をサポートしています。 |
|
username |
MySQL データベースサービスのユーザー名。 |
|
password |
MySQL データベースサービスのパスワード。 |
|
deserializer |
SourceRecord タイプのレコードを指定されたタイプに変換するデシリアライザ。有効な値は次のとおりです:
|
pom 依存関係
ローカルデバッグ
Realtime Compute for Apache Flink のコネクタには、商用化された追加コンテンツが含まれており、多くの点で Apache Flink とは異なります。ローカルデバッグを行うには、pom ファイルに以下の変更を加えます。
-
${flink.version} は 1.19.0 にする必要があります。
-
コネクタ ${vvr.version} の推奨バージョン:
-
VVR 8.x エンジンの場合、バージョン 1.17-vvr-8.0.11-4 を使用してください。
-
VVR 11.x エンジンの場合、お使いのエンジンバージョンに対応する最新のコネクタバージョンを使用してください。たとえば、VVR 11.8 エンジンの場合は、最新バージョンの 1.20-vvr-11.8.0-1-jdk11 を使用してください。その他のバージョンについては、「Maven リポジトリ」をご参照ください。
-
-
Kafka コネクタの依存関係を追加する必要があります。
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<java.version>8</java.version>
<flink.version>1.19.0</flink.version>
<vvr.version>1.20-vvr-11.8.0-1-jdk11</vvr.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-core</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<!-- VVR 11.x ではこの group ID を使用します -->
<groupId>com.alibaba.ververica</groupId>
<!-- VVR 8.x ではこの group ID を使用します -->
<!-- <groupId>org.apache.flink</groupId> -->
<artifactId>flink-table-common</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-kafka</artifactId>
<version>${vvr.version}</version>
</dependency>
</dependencies>
Realtime Compute for Apache Flink でのデプロイとデバッグ
Realtime Compute for Apache Flink でジョブをデバッグする場合、コネクタのバージョン制限は適用されません。${flink.version} と ${vvr.version} がジョブのエンジンバージョンに対応していることを確認してください。バージョンマッピングの詳細については、「エンジン」をご参照ください。以下の pom ファイルをご参照ください。
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-core</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-connector-base</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-common</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>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
</dependency>