すべてのプロダクト
Search
ドキュメントセンター

Data Transmission Service:flink-dts-connector を使用した追跡データの消費

最終更新日:Jun 21, 2026

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 を使用して変更追跡インスタンスから追跡データを消費する方法を示します。

  1. 変更追跡インスタンスを作成します。詳細については、「ApsaraDB RDS for MySQL インスタンスの変更追跡インスタンスを作成する」、「PolarDB for MySQL クラスターの変更追跡インスタンスを作成する」、または「ApsaraDB for Oracle インスタンスの変更追跡インスタンスを作成する」をご参照ください。

  2. 1 つ以上のコンシューマーグループを作成します。詳細については、「コンシューマーグループを作成する」をご参照ください。

  3. flink-dts-connector ファイルをダウンロードし、解凍します。

  4. IntelliJ IDEA を起動し、[Open or Import] をクリックします。

  5. 表示されたダイアログボックスで、flink-dts-connector ファイルを解凍したディレクトリに移動します。フォルダーを展開して、Project Object Model (POM) ファイル pom.xml を見つけます。

  6. 表示されるダイアログボックスで、[プロジェクトとして開く] を選択します。

  7. 次の依存関係を 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>
  8. IntelliJ IDEA でプロジェクトフォルダーを展開し、プログラムで使用する Flink API の種類に応じて、適切な Java ファイルを選択します。

    • Flink プログラムで DataStream API を使用する場合は、flink-dts-connector-master\src\test\java\com\alibaba\flink\connectors\dts\datastream\DtsExample.java ファイルをダブルクリックし、次の手順を実行します:

      1. 上部のメニューバーから、[Run] > [Run...] を選択します。

      2. ポップアップウィンドウで、[DtsExample] > [編集]をクリックします。

      3. [Program arguments] フィールドに、次の例のようにパラメーターとその値を入力します。次に、[Run] をクリックして flink-dts-connector を起動します。

        説明

        パラメーターの詳細と値の確認方法については、「パラメーター」をご参照ください。

        --broker-url dts-cn-******.******.***:****** --topic cn_hangzhou_rm_**********_dtstest_version2 --sid dts****** --user dtstest --password Test123456 --checkpoint 1624440043
      4. 出力は、プログラムがソースデータベースのデータ変更を正常に追跡していることを示します。起動後、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 つのテーブルからのみ追跡データを設定して消費できます。複数のテーブルからデータを消費するには、テーブルごとに個別のタスクを設定して実行する必要があります。

      1. // プレフィックスを追加して、properties.load 行をコメントアウトします。

        /* プラットフォームで設定したパラメーター値を Properties オブジェクトにロードします。 */
        //properties.load(new StringReader(new String(Files.readAllBytes(Paths.get(configFilePath)), StandardCharsets.UTF_8)));
      2. データを消費するテーブルを 1 つ定義します。

      3. 変更追跡インスタンスのパラメーターを設定します。パラメーターの詳細と値の確認方法については、「パラメーター」をご参照ください。

        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')";
        }
      4. IntelliJ IDEA の上部で、[Run'DtsTableISelectTCaseTest'] をクリックして flink-dts-connector を起動します。

      5. 出力は、プログラムがソースデータベースのデータ変更を正常に追跡していることを示します。起動後、ターミナルに変更ログレコードが表示されます。更新イベントはペアで表示されます。更新前の値 (-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 パラメーター

説明

取得方法

broker-url

dts.server

変更追跡インスタンスのエンドポイントとポートです。

説明
  • ネットワーク遅延を最小限に抑えるため、Flink プログラムの ECS インスタンスと変更追跡インスタンスが同じ仮想プライベートクラウド (VPC) またはクラシックネットワークにある場合は、内部エンドポイントを使用してください。

  • ネットワークが不安定になる可能性があるため、パブリックエンドポイントの使用は推奨されません。

Data Transmission Service (DTS) コンソールで対象の変更追跡インスタンスの ID をクリックします。[View Task Settings] ページで、トピック、エンドポイント、およびポートを確認できます。

topic

topic

変更追跡インスタンスの追跡対象トピックです。

sid

dts.sid

コンシューマーグループの ID です。

DTS コンソールで対象の変更追跡インスタンスの ID をクリックして、データ消費 をクリックします。コンシューマーグループの [ID] と アカウント を確認できます。

説明

コンシューマーグループアカウントのパスワードは、コンシューマーグループ作成時に指定したものです。

user

dts.user

コンシューマーグループのユーザー名です。

警告

このトピックで提供する flink-dts-connector は、必要なユーザー名の書式を自動的に処理します。別のクライアントを使用する場合は、接続失敗を防ぐために、ユーザー名を <username>-<consumer group ID> (例:dtstest-dtsae******bpv) の形式に手動で整形する必要があります。

password

dts.password

ユーザー名のパスワードです。

checkpoint

dts.checkpoint

コンシューマーオフセットです。これは、flink-dts-connector がデータの消費を開始するタイムスタンプです。値は Unix タイムスタンプである必要があります (例:1624440043)。

説明

コンシューマーオフセットは、次のシナリオで役立ちます:

  • コンシューマープログラムが中断された場合、最後に処理したオフセットを使用して消費を再開することで、データ損失を防ぐことができます。

  • コンシューマープログラムの起動時に、特定のオフセットを指定することで、目的の時点からデータの消費を開始できます。

コンシューマーオフセットは、変更追跡インスタンスのデータ範囲内に設定する必要があります。DTS 変更追跡インスタンスの [View Task Settings] ページの [Data Range] フィールドで、開始時刻と終了時刻を確認できます。この範囲内の時刻をコンシューマーオフセットとして設定し、Unix タイムスタンプに変換する必要があります。

説明

Unix タイムスタンプへの変換には、オンラインの変換ツールなどをご利用ください。

N/A

dts-cdc.table.name

このパラメーターで、変更追跡の対象テーブルを指定します。タスクあたり 1 つのテーブルのみサポートされます。次の形式を使用します:

  • [MySQL]、[PolarDB for MySQL]、[PolarDB-X 1.0]、および [PolarDB-X 2.0] データベースの場合は、<Database name>.<Table name> の形式を使用します。

  • その他のデータベースタイプの場合は、<Schema name>.<Table name> の形式を使用します。

DTS コンソールで、対象の変更追跡インスタンスの ID をクリックします。 [タスク設定の表示] ページで、右上隅の オブジェクトを表示 をクリックして、オブジェクトのデータベース名とテーブル名を確認します。

よくある質問

エラーメッセージ

考えられる原因

解決策

Cluster changed from *** to ***, consumer require restart.

DTS が増分データを読み取るために使用する DStore モジュールが切り替えられました。これにより、Flink プログラムのコンシューマーオフセットが失われます。

プログラムを再起動しないでください。代わりに、最後に把握しているコンシューマーオフセットを特定し、それを checkpoint または dts.checkpoint パラメーターとして DtsExample.java または DtsTableISelectTCaseTest.java ファイルに渡し、消費を再開してください。