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

Data Transmission Service:SDK サンプルコードを使用したデータサブスクリプションのデータ消費

最終更新日:Aug 06, 2026

データサブスクリプションインスタンスを設定した後、 Data Transmission Service (DTS) の SDK サンプルコードを使用して変更データを消費します。

操作手順

重要

次の手順では、IntelliJ IDEA (Community Edition 2020.1 for Windows) で SDK サンプルコードを実行して、データサブスクリプションのデータを消費する方法について説明します。

  1. データサブスクリプションインスタンスを作成します。詳細については、「ApsaraDB RDS for MySQL インスタンスのデータサブスクリプションチャネルの作成」、「PolarDB for MySQL クラスターのデータサブスクリプションチャネルの作成」、または「Oracle データベースのデータサブスクリプションチャネルの作成」をご参照ください。

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

    重要

    データサブスクリプションのデータを消費する場合、DefaultUserRecord の commit メソッドを呼び出してチェックポイントをコミットする必要があります。そうしないと、データが重複して消費される可能性があります。

  3. ビジネス要件に基づいて SDK サンプルコードを使用します。

    • 新しいデータサブスクリプション SDK パッケージを使用する (推奨)

      1. IntelliJ IDEA を開き、[Create New Project] をクリックして、アプリケーション用のプロジェクトを作成します。

      2. プロジェクトで、プロジェクトオブジェクトモデル (POM) ファイル pom.xml を見つけます。

      3. pom.xml ファイルに次の依存関係を追加します:

        <dependency>
            <groupId>com.aliyun.dts</groupId>
            <artifactId>dts-new-subscribe-sdk</artifactId>
            <version>{dts_new_sdk_version}</version>
        </dependency>
        説明

        最新の Maven 依存関係は、dts-new-subscribe-sdk のページで確認できます。

      4. 新しいサブスクリプション SDK の使用方法の詳細については、「サンプルコードの使用」をご参照ください。

    • 新しいデータサブスクリプション SDK のカスタマイズ版を使用する

      1. SDK サンプルコードパッケージをダウンロードし、解凍します。

        説明

        code をクリックし、[Download ZIP] を選択してパッケージをダウンロードします。

      2. 解凍した SDK サンプルコードのディレクトリに移動します。テキストエディターで pom.xml ファイルを開き、データサブスクリプション SDK を最新バージョンに更新します。

        <name>dts-new-subscribe-sdk</name>
        <url>https://www.aliyun.com/product/dts</url>
        <description>The Aliyun new Subscribe SDK for Java used for accessing Data Transmission Service</description>
        <packaging>jar</packaging>
        <groupId>com.aliyun.dts</groupId>
        <artifactId>dts-new-subscribe-sdk</artifactId>
        <version>1.3</version>
        重要

        データサブスクリプション SDK の最新バージョンは、Maven のウェブサイトから入手できます。詳細については、「データサブスクリプション SDK の Maven ページ」をご参照ください。

      3. IntelliJ IDEA を開きます。ようこそ画面で、[Open or Import] をクリックします。

      4. [Open File or Project] ダイアログボックスで、解凍した SDK サンプルコードのディレクトリに移動し、pom.xml ファイルを選択して [OK] をクリックします。

      5. 表示されるダイアログボックスで、[Open as Project] を選択します。

      6. IntelliJ IDEA の [Project] パネルで、src > test > java > com.aliyun.dts.subscribe.clients を展開します。[DTSConsumerAssignDemo] と [DTSConsumerSubscribeDemo] は、2 つのサンプルデモファイルです。SDK クライアントの使用モードに基づいて、対応する Java ファイル (DTSConsumerAssignDemo.java または DTSConsumerSubscribeDemo.java) を選択し、ダブルクリックします。

        プロジェクトファイルツリー:

        aliyun-dts-subscribe-sdk-java-master [dts-new-subscribe-sdk] D:\aliyun-dts-sub
          .idea
          src
            main
            test
              java
                com.aliyun.dts.subscribe.clients
                  DTSConsumerAssignDemo
                  DTSConsumerSubscribeDemo
                  UserMetaStore
          target
          .gitignore
          dts-new-subscribe-sdk.iml
          LICENSE
          pom.xml
          README.md
          External Libraries
          Scratches and Consoles
        説明

        DTS は、次の SDK クライアントの使用モードをサポートしています:

        • ASSIGN モード:メッセージのグローバルな順序を保証するために、DTS は各サブスクリプショントピックに 1 つのパーティション (パーティション 0) のみを割り当てます。SDK クライアントを ASSIGN モードで使用する場合、1 つのクライアントのみを起動することを推奨します。

        • SUBSCRIBE モード:メッセージのグローバルな順序を保証するために、DTS は各サブスクリプショントピックに 1 つのパーティション (パーティション 0) のみを割り当てます。SDK クライアントを SUBSCRIBE モードで使用する場合、災害復旧のためにコンシューマーグループ内で複数の SDK クライアントを起動できます。アクティブなクライアントに障害が発生した場合、別の SDK クライアントが自動的にパーティション 0 に割り当てられ、消費を再開します。

  4. Java ファイルで必要なパラメーターを設定します。

    public static void main(String[] args) {
        // Kafka ブローカーの URL
        String brokerUrl = "dts-cn-xxx18001";
        // 消費するトピック、パーティションは 0
        String topic = "cn_hangzhou_rm_xxx";
        // 認証用のユーザーパスワードと sid
        String sid = "dtsxxx";
        String userName = "dtstest";
        String password = "xxx";
        // 最初のシークのための初期チェックポイント (設定するタイムスタンプ、例:1566180200 (2019 年 8 月 19 日月曜日 10:03:21 CST の場合))
        String initCheckpoint = "1620811813";
        // SUBSCRIBE モードを使用する場合、グループ設定が必要です。Kafka コンシューマーグループが有効になります
        ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.ASSIGN;
        // 開始時に強制的に設定チェックポイントを使用するかどうか。チェックポイントのリセットは、ASSIGN モードでのみ機能します
        boolean isForceUseInitCheckpoint = true;

    表 1. 必須パラメーター

    パラメーター

    説明

    ソース

    brokerUrl

    データサブスクリプションインスタンスのエンドポイントとポート番号。

    説明
    • SDK クライアントを実行している ECS インスタンスとデータサブスクリプションインスタンスが同じクラシックネットワークまたは Virtual Private Cloud (VPC) にある場合は、ネットワークレイテンシを最小限に抑えるために、サブスクリプションに内部エンドポイントを使用することを推奨します。

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

    DTS コンソールで、対象のデータサブスクリプションインスタンスの ID をクリックします。 基本情報 ページで、ネットワーク セクションのエンドポイントとポート番号を取得できます。

    topic

    インスタンスのサブスクリプショントピック。

    DTS コンソールで、対象のデータサブスクリプションインスタンスの ID をクリックします。 基本情報 ページの 基本情報 セクションにある トピック を取得できます。

    sid

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

    DTS コンソールで、対象のデータサブスクリプションインスタンスの ID をクリックし、データ消費 をクリックします。コンシューマーグループの [コンシューマーグループ ID] と アカウント を取得できます。

    説明

    コンシューマーグループのユーザー名のパスワードは、コンシューマーグループを作成するときに指定します。

    userName

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

    警告

    このトピックで提供されているクライアントを使用しない場合、ユーザー名を <Username>-<Consumer Group ID> 形式で設定する必要があります。 例: dtstest-dtsae******bpv。 そうしないと、接続に失敗します。

    password

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

    initCheckpoint

    SDK クライアントがデータの消費を開始する消費チェックポイント。UNIX タイムスタンプとして指定します。例:1620962769。

    説明

    次のシナリオで消費チェックポイント情報を使用できます:

    • アプリケーションが中断した後に消費を再開し、データの損失を防ぐには、最後に認識された消費チェックポイントを渡してください。

    • クライアントを起動するときに、特定の消費チェックポイントを渡して、目的の位置からデータを消費できます。

    消費チェックポイントは、データサブスクリプションインスタンスのデータ範囲 (図を参照) 内にある必要があり、UNIX タイムスタンプに変換する必要があります。

    説明

    検索エンジンを使用して、UNIX タイムスタンプコンバーターを見つけることができます。

    ConsumerContext.ConsumerSubscribeMode subscribeMode

    SDK クライアントの使用モード。有効な値:

    • ConsumerContext.ConsumerSubscribeMode.ASSIGN:ASSIGN モード。コンシューマーグループ内の 1 つの SDK クライアントのみがデータサブスクリプションのデータを消費できます。

    • ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE:SUBSCRIBE モード。災害復旧のために、同じコンシューマーグループで複数の SDK クライアントを起動できます。

    N/A

  5. IntelliJ IDEA の上部メニューで、[Run] > [Run] を選択してクライアントを実行します。

    説明

    クライアントを初めて実行するときは、必要な依存関係の自動ロードとインストールに時間がかかる場合があります。

    • 実行後、SDK クライアントはソースデータベースから変更データを正常に消費します。以下は、消費された UPDATE レコードの例です:

      [2021-05-18 16:49:50,260] INFO
      RecordID [559686]
      RecordTimestamp [1621327772]
      Source [{"sourceType": "MySQL", "version": "8.0.18"}]
      RecordType [UPDATE]
      Schema info [{,
      recordFields= [{fieldName='orderid', rawDataTypeNum=3, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}, {fieldName='username', rawDataTypeNum=254, isPrimaryKey=false, isUniqueKey=false, fieldPosition=1}, {fieldName='ordertime', rawDataTypeNum=12, isPrimaryKey=false, isUniqueKey=false, fieldPosition=2}, {fieldName='commodity', rawDataTypeNum=253, isPrimaryKey=false, isUniqueKey=false, fieldPosition=3}, {fieldName='phonenumber', rawDataTypeNum=3, isPrimaryKey=false, isUniqueKey=false, fieldPosition=4}, {fieldName='address', rawDataTypeNum=15, isPrimaryKey=false, isUniqueKey=false, fieldPosition=5}],
      databaseName='dtstestdata',
      tableName='order',
      primaryIndexInfo [{indexType=PrimaryKey, indexFields=[{fieldName='orderid', rawDataTypeNum=3, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}], cardinality=0, nullable=true, isFirstUniqueIndex=false, name=null}],
      uniqueIndexInfo [[]],
      partitionFields = null}]
      変更前イメージ {[Field [orderid] [1]
      Field [username] [ 田中花子]
      Field [ordertime] [xxx]
      Field [commodity] [xxx]
      Field [phonenumber] [ xxx]
      Field [address] [xxx]
      ]}
      変更後イメージ {[Field [orderid] [1]
      Field [username] [ 花子]
      Field [ordertime] [xxx]
      Field [commodity] [xxx]
      Field [phonenumber] [ xxx]
      Field [address] [xxx]
      ]}
        (com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener)
    • SDK クライアントは、送受信されたレコードの総数と量、および RPS (Requests Per Second) を含むデータ消費に関する統計を定期的に集計して表示します。

      [2021-05-18 16:25:09,167] INFO {"outCounts":488616.0,"outBytes":48606134,"outRps":1.15,"outBps":114.57,"count":11.0,"inBytes":60118961,"DStoreRecordQueue":0.0,"inCounts":557154.0,"inRps":1.12,"inBps":112.44,"__dt":1621326309167,"DefaultUserRecordQueue":0.0} (log_metrics)

      表 2. データ消費統計

      パラメーター

      説明

      outCounts

      SDK クライアントによって消費されたデータレコードの総数。

      outBytes

      SDK クライアントによって消費されたデータの総量。単位:バイト。

      outRps

      SDK クライアントがデータを消費する際の RPS (Requests Per Second)。

      outBps

      SDK クライアントのデータ消費レート。単位:バイト/秒。

      inBytes

      DTS サーバーから送信されたデータの総量。単位:バイト。

      DStoreRecordQueue

      DTS サーバーから受信したレコードを格納する内部データキャッシュキューのサイズ。

      inCounts

      DTS サーバーから送信されたデータレコードの総数。

      inRps

      DTS サーバーがデータを送信する際の RPS。

      __dt

      SDK クライアントがデータを受信したときのタイムスタンプ。単位:ミリ秒。

      DefaultUserRecordQueue

      コンシューマーアプリケーションが処理する準備ができたレコードを保持するデータキューのサイズ。

消費チェックポイントの保存とクエリ

データ消費を開始または再開する際 (初回起動時、再起動時、内部リトライ時など)、SDK クライアントには消費チェックポイントが必要です。次の表では、データ損失を防ぎ、重複消費を最小限に抑え、オンデマンド消費を可能にするために、さまざまなシナリオでチェックポイントを管理およびクエリする方法について説明します。

シナリオ

SDK 利用モード

クエリ方法

消費チェックポイントのクエリ

ASSIGN モード、SUBSCRIBE モード

  • SDK クライアントは 5 秒ごとに消費チェックポイントを保存し、DTS サーバーに送信します。最新の消費チェックポイントをクエリするには、次のいずれかの方法を使用します。

    • SDK クライアントが実行されているサーバー上の localCheckpointStore ファイルを確認してください。

    • データサブスクリプションインスタンスの データ消費 ページを確認してください。

  • consumerContext.java ファイルで setUserRegisteredStore(newUserMetaStore()) を使用して外部の永続的な共有ストレージメディア (データベースなど) を設定している場合、このストレージメディアは 5 秒ごとに消費チェックポイントを保存するため、クエリできます。

初回起動: チェックポイントを渡して消費を開始する。

ASSIGN モード、SUBSCRIBE モード

SDK クライアントの利用モードに基づいて、DTSConsumerAssignDemo.java または DTSConsumerSubscribeDemo.java ファイルを選択し、initCheckpoint パラメーターを設定して消費を開始してください。設定手順については、手順 3 および 4 をご参照ください。

内部リトライ後に消費を継続するため、SDK クライアントは最後に記録された消費チェックポイントを再度渡す必要がある。

ASSIGN モード

次の順序で最後に記録された消費チェックポイントを検索します。検索は、チェックポイント情報が見つかり次第停止し、その情報を返します。

  1. consumerContext.java ファイルで setUserRegisteredStore(newUserMetaStore()) を使用して設定した外部ストレージメディア。

  2. SDK クライアントが実行されているサーバー上の localCheckpointStore ファイル。

  3. DTSConsumerAssignDemo.java ファイルの initCheckpoint パラメーターで渡した開始タイムスタンプ。

SUBSCRIBE モード

次の順序で最後に記録された消費チェックポイントを検索します。検索は、チェックポイント情報が見つかり次第停止し、その情報を返します。

  1. consumerContext.java ファイルで setUserRegisteredStore(newUserMetaStore()) を使用して設定した外部ストレージメディア。

  2. DTS サーバー (増分データ収集モジュール) に保存されたチェックポイント。

    説明

    このオフセットは、SDK クライアントが commit メソッドを呼び出してコンシューマオフセットを更新した後にのみ更新されます。

  3. DTSConsumerSubscribeDemo.java ファイルの initCheckpoint パラメーターで渡した開始タイムスタンプ。

  4. DTS サーバー (新しい増分データ収集モジュール) の開始チェックポイント。

    重要

    増分データ収集モジュールの切り替えが発生した場合、新しいモジュールはクライアントからの最後の消費チェックポイントを保持しません。これにより、クライアントが古いチェックポイントからデータの消費を開始する可能性があります。クライアント側で消費チェックポイントを永続化することを推奨します。詳細については、「消費チェックポイントの永続化」をご参照ください。

SDK クライアントが再起動され、消費を継続するために最後に記録された消費チェックポイントを再度渡す必要がある。

ASSIGN モード

consumerContext.java ファイルの setForceUseCheckpoint 設定に基づいて消費チェックポイントをクエリします。検索は、チェックポイント情報が見つかり次第停止し、その情報を返します。

  • true に設定されている場合、SDK クライアントの再起動時に毎回、渡された initCheckpoint が消費チェックポイントとして強制的に使用されます。

  • false に設定されているか、設定されていない場合は、次の順序で最後に記録された消費チェックポイントを検索します。

    1. SDK クライアントが実行されるサーバー上の localCheckpointStore ファイル。

    2. DTS サーバー (増分データ収集モジュール) に保存されたチェックポイント。

      説明

      このオフセットは、SDK クライアントが commit メソッドを呼び出してコンシューマオフセットを更新した後にのみ更新されます。

    3. consumerContext.java ファイルで setUserRegisteredStore(newUserMetaStore()) を使用して設定した外部ストレージメディア。

SUBSCRIBE モード

このモードでは、consumerContext.java ファイルの setForceUseCheckpoint 設定は無効です。次の順序で最後に記録された消費チェックポイントを検索します。

  1. consumerContext.java ファイルで setUserRegisteredStore(newUserMetaStore()) を使用して設定した外部ストレージメディア。

  2. DTS サーバー (増分データ収集モジュール) に保存されたチェックポイント。

    説明

    このオフセットは、SDK クライアントが commit メソッドを呼び出してコンシューマオフセットを更新した後にのみ更新されます。

  3. DTSConsumerSubscribeDemo.java ファイルの initCheckpoint パラメーターで渡した開始タイムスタンプ。

  4. DTS サーバー (新しい増分データ収集モジュール) の開始チェックポイント。

消費チェックポイントの永続化

増分データ収集モジュールでディザスタリカバリイベントが発生した場合 (特に SUBSCRIBE モードにおいて)、新しいモジュールはクライアントの最新の消費チェックポイントを保持しません。クライアントは古いチェックポイントから再開する可能性があり、履歴データの重複消費が発生します。たとえば、切り替え前に、旧モジュールのチェックポイント範囲が 2023 年 11 月 11 日 08:00:00 から 2023 年 11 月 12 日 08:00:00 で、クライアントのチェックポイントが 2023 年 11 月 12 日 08:00:00 である場合、切り替え後、新しいモジュールのチェックポイント範囲は 2023 年 11 月 08 日 10:00:00 から 2023 年 11 月 12 日 08:01:00 になります。クライアントは新しいモジュールの開始チェックポイント (2023 年 11 月 08 日 10:00:00) から開始するため、重複消費が発生します。

このシナリオで重複消費を回避するには、クライアントに永続的なチェックポイントストアを設定してください。次の例では、要件に応じて調整できる実装例を 1 つ示します。

  1. AbstractUserMetaStore を extends する UserMetaStore クラスを作成してください。

    たとえば、チェックポイント情報を MySQL データベースに保存するには、次の Java コードを使用してください。

    public class UserMetaStore extends AbstractUserMetaStore {
    
        @Override
        protected void saveData(String groupID, String toStoreJson) {
            Connection con = getConnection();
            String sql = "insert into dts_checkpoint(group_id, checkpoint) values(?, ?)";
    
            PreparedStatement pres = null;
            ResultSet rs = null;
    
            try {
                pres = con.prepareStatement(sql);
                pres.setString(1, groupID);
                pres.setString(2, toStoreJson);
                pres.execute();
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                close(rs, pres, con);
            }
        }
    
        @Override
        protected String getData(String groupID) {
            Connection con = getConnection();
            String sql = "select checkpoint from dts_checkpoint where group_id = ?";
    
            PreparedStatement pres = null;
            ResultSet rs = null;
    
            try {
                pres = con.prepareStatement(sql);
                pres.setString(1, groupID);
                rs = pres.executeQuery();
                              
                String checkpoint = rs.getString("checkpoint");
              
                return checkpoint;
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                close(rs, pres, con);
            }
        }
    }
    
  2. consumerContext.java ファイルで、setUserRegisteredStore(new UserMetaStore()) メソッドを使用して外部ストレージを設定してください。

よくある質問

  • データサブスクリプションインスタンスとの接続に関する問題を解決するには、どうすればよいですか。

    エラーメッセージに基づいてトラブルシューティングを行ってください。詳細については、「トラブルシューティング」をご参照ください。

  • 消費チェックポイントはどのような形式で永続化されますか。

    永続化された消費チェックポイントデータは、JSON フォーマットで保存されます。この永続化されたチェックポイントは、SDK に直接渡すことができる UNIX タイムスタンプです。次のレスポンス例では、"timestamp" キーの値 1700709977 が永続化された消費チェックポイントです。

    {"groupID":"dtsglg11d48230***","streamCheckpoint":[{"partition":0,"offset":577989,"topic":"ap_southeast_1_vpc_rm_t4n22s21iysr6****_root_version2","timestamp":1700709977,"info":""}]}

トラブルシューティング

問題

エラーメッセージ

原因

解決策

接続できません

ERROR
CheckResult{isOk=false, errMsg='telnet dts-cn-hangzhou.aliyuncs.com:18009
failed, please check the network and if the brokerUrl is correct'}
(com.aliyun.dts.subscribe.clients.DefaultDTSConsumer)

指定された brokerUrl が正しくありません。

brokerUrl、userName、および password パラメーターに正しい値を入力します。詳細については、「必須パラメーター」をご参照ください。

telnet real node *** failed, please check the network

ブローカーアドレスが実際の IP アドレスに接続できません。

ERROR CheckResult{isOk=false, errMsg='build kafka consumer failed, error: org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata, probably the user name or password is wrong'} (com.aliyun.dts.subscribe.clients.DefaultDTSConsumer)

ユーザー名またはパスワードが正しくありません。

com.aliyun.dts.subscribe.clients.exception.TimestampSeekException: RecordGenerator:seek timestamp for topic [cn_hangzhou_rm_bp11tv2923n87081s_rdsdt_dtsacct-0] with timestamp [1610249501] failed

consumerContext.java ファイルで、setUseCheckpoint パラメーターが true に設定されていますが、消費チェックポイントがデータサブスクリプションインスタンスのデータ範囲内にありません (図を参照)。

データサブスクリプションインスタンスのデータ範囲内にある消費チェックポイントを渡します。詳細については、「必須パラメーター」をご参照ください。

消費速度の低下

N/A

  • 統計で DStoreRecordQueue と DefaultUserRecordQueue のサイズを確認し、ボトルネックを特定します。詳細については、「データ消費統計」をご参照ください。

    • DStoreRecordQueue が 0 のままである場合、DTS サーバーが低速でデータをプルしていることを示します。

    • DefaultUserRecordQueue が一貫してそのキャパシティ (デフォルトは 512) に近い状態が続く場合、SDK クライアントがデータを消費するのが遅すぎることを示します。

  • ビジネス要件に基づいて、コード内の消費チェックポイント (initCheckpoint) を変更して、チェックポイントをリセットします。