Data Transmission Service (DTS) の変更追跡インスタンスを設定した後、flink-dts-connector を使用して、追跡データを消費する Flink プログラムを作成できます。
注意事項
-
このコネクタは、DataStream API、または Table API と SQL を使用する Flink プログラムをサポートします。
-
Flink プログラムで Table API と SQL を使用する場合、同時に 1 つのテーブルからのみデータを消費するように設定できます。複数のテーブルからデータを消費するには、テーブルごとに個別のタスクを設定して実行する必要があります。
手順
このトピックでは、Windows 版 IntelliJ IDEA Community Edition 2020.1 を例に、flink-dts-connector を使用して変更追跡インスタンスから追跡データを消費する方法を示します。
-
変更追跡インスタンスを作成します。詳細については、「ApsaraDB RDS for MySQL インスタンスの変更追跡インスタンスを作成する」、「PolarDB for MySQL クラスターの変更追跡インスタンスを作成する」、または「ApsaraDB for Oracle インスタンスの変更追跡インスタンスを作成する」をご参照ください。
-
1 つ以上のコンシューマーグループを作成します。詳細については、「コンシューマーグループを作成する」をご参照ください。
-
flink-dts-connector ファイルをダウンロードし、解凍します。
-
IntelliJ IDEA を起動し、[Open or Import] をクリックします。
-
表示されたダイアログボックスで、flink-dts-connector ファイルを解凍したディレクトリに移動します。フォルダーを展開して、Project Object Model (POM) ファイル
pom.xmlを見つけます。 表示されるダイアログボックスで、[プロジェクトとして開く] を選択します。
-
次の依存関係を
pom.xmlファイルに追加します:<dependency> <groupId>com.alibaba.flink</groupId> <artifactId>flink-dts-connector</artifactId> <version>1.1.1-SNAPSHOT</version> <classifier>jar-with-dependencies</classifier> </dependency> -
IntelliJ IDEA でプロジェクトフォルダーを展開し、プログラムで使用する Flink API の種類に応じて、適切な Java ファイルを選択します。
-
Flink プログラムで DataStream API を使用する場合は、
flink-dts-connector-master\src\test\java\com\alibaba\flink\connectors\dts\datastream\DtsExample.javaファイルをダブルクリックし、次の手順を実行します:-
上部のメニューバーから、[Run] > [Run...] を選択します。
-
ポップアップウィンドウで、をクリックします。
-
[Program arguments] フィールドに、次の例のようにパラメーターとその値を入力します。次に、[Run] をクリックして flink-dts-connector を起動します。
説明パラメーターの詳細と値の確認方法については、「パラメーター」をご参照ください。
--broker-url dts-cn-******.******.***:****** --topic cn_hangzhou_rm_**********_dtstest_version2 --sid dts****** --user dtstest --password Test123456 --checkpoint 1624440043 -
出力は、プログラムがソースデータベースのデータ変更を正常に追跡していることを示します。起動後、DataStream 内のデータ変更レコードは次の例のようになります。
LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006303@289540@2688@1625045211000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006307@290047@2688@1625045212000]} LazyParseRecord {operationType [UPDATE], checkpoint [0@12006305@290016@2688@1625045212000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006308@290047@2688@1625045214000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006309@290047@2688@1625045215000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006310@290047@2688@1625045216000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006311@290047@2688@1625045217000]}説明データ変更レコードの詳細を表示するには、Flink プログラムの TaskManager UI にログインします。
-
-
Flink プログラムで Table API と SQL を使用する場合は、
flink-dts-connector-master\src\test\java\com\alibaba\flink\connectors\dts\sql\DtsTableISelectTCaseTest.javaファイルをダブルクリックし、次の手順を実行します:説明DtsTableISelectTCaseTest.javaファイルでは、1 つのテーブルからのみ追跡データを設定して消費できます。複数のテーブルからデータを消費するには、テーブルごとに個別のタスクを設定して実行する必要があります。-
//プレフィックスを追加して、properties.load行をコメントアウトします。/* プラットフォームで設定したパラメーター値をPropertiesオブジェクトにロードします。 */ //properties.load(new StringReader(new String(Files.readAllBytes(Paths.get(configFilePath)), StandardCharsets.UTF_8))); -
データを消費するテーブルを 1 つ定義します。
-
変更追跡インスタンスのパラメーターを設定します。パラメーターの詳細と値の確認方法については、「パラメーター」をご参照ください。
public static void main(String[] args) throws Exception { setup(args); final String createTable = "create table `dts` (\n" + " `ts` TIMESTAMP(3) METADATA FROM 'timestamp',\n" + " `id` bigint,\n" + " `name` varchar,\n" + " `age` bigint,\n" + " WATERMARK FOR ts AS ts - INTERVAL '5' SECOND" + ") with (\n" + "'connector' = 'dts'," + "'dts.server' = 'dts-cn-xxx:18001'," + "'topic' = 'cn_hangzhou_rm_xxx_dtstest_version2'," + "'dts.sid' = 'dtsxxx', " + "'dts.user' = 'dtstest', " + "'dts.password' = 'xxx'," + "'dts.checkpoint' = '1624440043', " + "'dts-cdc.table.name' = 'dtstestdata.order'," + "'format' = 'dts-cdc')"; } -
IntelliJ IDEA の上部で、[Run'DtsTableISelectTCaseTest'] をクリックして flink-dts-connector を起動します。
-
出力は、プログラムがソースデータベースのデータ変更を正常に追跡していることを示します。起動後、ターミナルに変更ログレコードが表示されます。更新イベントはペアで表示されます。更新前の値 (-U) と更新後の値 (+U) です。
######> (false,-U(2021-06-23T20:32:17.391,null,null,null)) ######> (true,+U(2021-06-23T20:32:17.391,null,null,null)) ######> (false,-U(2021-06-23T20:32:45.604,null,null,null)) ######> (true,+U(2021-06-23T20:32:45.604,null,null,null)) ######> (false,-U(2021-06-30T17:26:52.201,null,null,null)) ######> (true,+U(2021-06-30T17:26:52.201,null,null,null)) ######> (false,-U(2021-06-30T19:19:26.975,null,null,null)) ######> (true,+U(2021-06-30T19:19:26.975,null,null,null))説明データ変更レコードの詳細を表示するには、Flink プログラムの TaskManager UI にログインします。
-
-
パラメーター
|
DataStream API パラメーター |
Table API パラメーター |
説明 |
取得方法 |
|
|
|
変更追跡インスタンスのエンドポイントとポートです。 説明
|
Data Transmission Service (DTS) コンソールで対象の変更追跡インスタンスの ID をクリックします。[View Task Settings] ページで、トピック、エンドポイント、およびポートを確認できます。 |
|
|
|
変更追跡インスタンスの追跡対象トピックです。 |
|
|
|
|
コンシューマーグループの ID です。 |
DTS コンソールで対象の変更追跡インスタンスの ID をクリックして、データ消費 をクリックします。コンシューマーグループの [ID] と アカウント を確認できます。 説明
コンシューマーグループアカウントのパスワードは、コンシューマーグループ作成時に指定したものです。 |
|
|
|
コンシューマーグループのユーザー名です。 警告
このトピックで提供する flink-dts-connector は、必要なユーザー名の書式を自動的に処理します。別のクライアントを使用する場合は、接続失敗を防ぐために、ユーザー名を |
|
|
|
|
ユーザー名のパスワードです。 |
|
|
|
|
コンシューマーオフセットです。これは、flink-dts-connector がデータの消費を開始するタイムスタンプです。値は Unix タイムスタンプである必要があります (例:1624440043)。 説明
コンシューマーオフセットは、次のシナリオで役立ちます:
|
コンシューマーオフセットは、変更追跡インスタンスのデータ範囲内に設定する必要があります。DTS 変更追跡インスタンスの [View Task Settings] ページの [Data Range] フィールドで、開始時刻と終了時刻を確認できます。この範囲内の時刻をコンシューマーオフセットとして設定し、Unix タイムスタンプに変換する必要があります。 説明
Unix タイムスタンプへの変換には、オンラインの変換ツールなどをご利用ください。 |
|
N/A |
|
このパラメーターで、変更追跡の対象テーブルを指定します。タスクあたり 1 つのテーブルのみサポートされます。次の形式を使用します:
|
DTS コンソールで、対象の変更追跡インスタンスの ID をクリックします。 [タスク設定の表示] ページで、右上隅の オブジェクトを表示 をクリックして、オブジェクトのデータベース名とテーブル名を確認します。 |
よくある質問
|
エラーメッセージ |
考えられる原因 |
解決策 |
|
DTS が増分データを読み取るために使用する DStore モジュールが切り替えられました。これにより、Flink プログラムのコンシューマーオフセットが失われます。 |
プログラムを再起動しないでください。代わりに、最後に把握しているコンシューマーオフセットを特定し、それを |