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

Data Transmission Service:ソフトウェア開発キット (SDK) を使用したサブスクライブ済みデータの消費

最終更新日:Sep 15, 2026

変更追跡タスクとコンシューマーグループを作成して変更追跡チャネルを構成した後、Data Transmission Service (DTS) が提供するソフトウェア開発キット (SDK) を使用して、サブスクライブしたデータを消費できます。このトピックでは、サンプルコードの使用方法について説明します。

説明

前提条件

注意事項

  • サブスクライブしたデータを消費する際は、データ消費を完了してから DefaultUserRecord の commit メソッドを呼び出してオフセット情報をコミットする必要があります。データを消費する前に commit メソッドを呼び出さないでください。そうしないと、データ損失が発生します。

  • 異なる消費プロセスは互いに独立しています。

  • コンソールでは、[現在のオフセット] は、追跡タスクがサブスクライブしたオフセットを示しており、クライアントによってコミットされたオフセットではありません。

操作手順

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

  2. SDK コードのバージョンを確認します。

    1. サンプル SDK コードを解凍したディレクトリに移動します。

    2. テキストエディタを使用して、ディレクトリ内の pom.xml ファイルを開きます。

    3. 変更追跡 SDK を最新バージョンに更新します。

      説明

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

      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>2.1.4</version>
  3. SDK コードを編集します。

    1. 統合開発環境 (IDE) を使用して、解凍したファイルを開きます。

    2. SDK クライアントを使用するモードに対応する Java ファイルを開きます。

      説明

      Java ファイルのパスは aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/ です。

      使用モード

      Java ファイル

      説明

      シナリオ

      ASSIGN モード

      DTSConsumerAssignDemo.java

      グローバルメッセージの順序を保証するため、DTS は各変更追跡トピックに 1 つのパーティション (パーティション 0) のみを割り当てます。SDK クライアントを ASSIGN モードで使用する場合は、SDK クライアントは 1 つだけ起動してください。

      コンシューマーグループ内の 1 つの SDK クライアントのみがサブスクライブしたデータを消費します。

      SUBSCRIBE モード

      DTSConsumerSubscribeDemo.java

      グローバルメッセージの順序を保証するため、DTS は各変更追跡トピックに 1 つのパーティション (パーティション 0) のみを割り当てます。SDK クライアントを SUBSCRIBE モードで使用する場合、同じコンシューマーグループ内で複数の SDK クライアントをディザスタリカバリ目的で起動できます。データを消費しているクライアントに障害が発生した場合、別の SDK クライアントがランダムかつ自動的にパーティション 0 に割り当てられ、消費を継続します。

      同じコンシューマーグループ内の複数の SDK クライアントがサブスクライブしたデータを消費します。これはディザスタリカバリシナリオです。

    3. Java コード内のパラメータを設定します。

      サンプルコード

      ******        
          public static void main(String[] args) {
              // Kafka ブローカーの URL
              String brokerUrl = "dts-cn-***.com:18001";
              // データを消費するトピック。パーティションは 0
              String topic = "cn_***_version2";
              // 認証用のユーザー名、パスワード、SID
              String sid = "dts***";
              String userName = "dts***";
              String password = "DTS***";
              // 最初のシークの初期チェックポイント。これは UNIX タイムスタンプです。たとえば、2019 年 8 月 19 日月曜日 10:03:20 (CST) から消費を開始する場合は、このパラメータを 1566180200 に設定します
              String initCheckpoint = "1740472***";
              // SUBSCRIBE モードを使用する場合は、グループを構成する必要があります。Kafka コンシューマーグループが有効になります
              ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE;
        
              DTSConsumerSubscribeDemo consumerDemo = new DTSConsumerSubscribeDemo(brokerUrl, topic, sid, userName, password, initCheckpoint, subscribeMode);
              consumerDemo.start();
          }
      ******

      パラメータ

      説明

      取得方法

      brokerUrl

      変更追跡チャネルのネットワークアドレスとポート番号を指定します。

      説明
      • Elastic Compute Service (ECS) インスタンスなど、SDK クライアントをデプロイするサーバーと変更追跡インスタンスが同じ Virtual Private Cloud (VPC) 内にある場合、VPC 経由でデータを消費することでネットワーク遅延を削減できます。

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

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

      topic

      変更追跡チャネルのトピックです。

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

      sid

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

      DTS コンソールで、ターゲットサブスクリプションインスタンス ID をクリックします。データ消費 ページで、コンシューマーグループ ID /名前 と アカウント を取得します。

      userName

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

      警告

      このドキュメントで提供されているクライアントを使用しない場合は、ユーザー名を <コンシューマーグループユーザー名>-<コンシューマーグループ ID> の形式で設定してください (例: dtstest-dtsae******bpv)。そうしないと、接続が失敗します。

      password

      アカウントのパスワードです。

      コンシューマーグループの作成時に、コンシューマーグループユーザー名に対して設定したパスワードです。

      initCheckpoint

      コンシューマーオフセットです。これは、SDK クライアントが消費する最初のデータレコードの UNIX タイムスタンプです (例: 1620962769)。

      説明

      コンシューマーオフセットを使用して、以下のことができます:

      • アプリケーションが中断された後、特定のオフセットからデータ消費を再開し、データ損失を防ぐことができます。

      • 必要に応じてデータ消費の開始オフセットを調整できます。

      コンシューマーオフセットは、変更追跡インスタンスのタイムスタンプ範囲内である必要があり、UNIX タイムスタンプに変換する必要があります。

      説明

      追跡タスク一覧にあるデータ範囲列には、ターゲットサブスクリプションインスタンスのタイムスタンプの範囲が表示されます。

      subscribeMode

      SDK クライアントの使用モードです。このパラメータを変更する必要はありません。

      • ConsumerContext.ConsumerSubscribeMode.ASSIGN: ASSIGN モード

      • ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE: SUBSCRIBE モード

      N/A

  4. IDE でプロジェクト構造を開き、プロジェクトの OpenJDK バージョンが 1.8 であることを確認します。

  5. クライアントコードを実行します。

    説明

    初回実行時、IDE は必要なプラグインと依存関係を自動的にロードするため、時間がかかることがあります。

    サンプル結果(クリックして展開)

    正常な実行結果

    以下の結果が返された場合、クライアントは正常に実行されており、ソースデータベースからのデータ変更をサブスクライブできる状態です。

    ******
    [2025-02-25 18:47:22.991] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.consumer.KafkaConsumer:1587] - [Consumer clientId=consumer-dtsl5vy2ao5250****-1, groupId=dtsl5vy2ao5250****] Seeking to offset 8200 for partition cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0
    [2025-02-25 18:47:22.993] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ConsumerWrap:116] - RecordFetcher consumer:  subscribe for [cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0] with checkpoint [Checkpoint[ topicPartition: cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0timestamp: 174048****, offset: 8200, info: 174048****]] start
    [2025-02-25 18:47:23.011] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174048044****,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:47:23.226] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8200]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [HEARTBEAT]
    
    [2025-02-25 18:47:23.226] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8201]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [HEARTBEAT]
    ******

    正常なサブスクリプション結果

    以下の結果が返された場合、クライアントはソースデータベースからのデータ変更 (UPDATE 操作) を正常にサブスクライブしています。

    ******
    [2025-02-25 18:48:24.905] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8413]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [UPDATE]
    Schema info [{, 
    recordFields= [{fieldName='id', rawDataTypeNum=8, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}, {fieldName='name', rawDataTypeNum=253, isPrimaryKey=false, isUniqueKey=false, fieldPosition=1}], 
    databaseName='dtsdb', 
    tableName='person', 
    primaryIndexInfo [[indexType=PrimaryKey, indexFields=[{fieldName='id', rawDataTypeNum=8, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}], cardinality=0, nullable=true, isFirstUniqueIndex=false, name=null]], 
    uniqueIndexInfo [[]], 
    partitionFields = null}]
    Before image {[Field [id] [3]
    Field [name] [test1]
    ]}
    After image {[Field [id] [3]
    Field [name] [test2]
    ]}
    ******

    異常な実行結果

    以下の結果が返された場合、クライアントはソースデータベースに接続できません。

    ******
    [2025-02-25 18:22:18.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174047893****,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:22:22.002] [WARN ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.NetworkClient:780] - [Consumer clientId=consumer-dtsnd7u2n0625m****-1, groupId=dtsnd7u2n0625m****] Connection to node 1 (47.118.XXX.XXX/47.118.XXX.XXX:18001) could not be established. Broker may not be available.
    [2025-02-25 18:22:22.509] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ClusterSwitchListener:44] - Cluster not changed on update:5aPLLlDtTHqP8sKq-DZVfg
    [2025-02-25 18:22:23.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":1740478943160,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:22:27.192] [WARN ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.NetworkClient:780] - [Consumer clientId=consumer-dtsnd7u2n0625m****1, groupId=dtsnd7u2n0625m****] Connection to node 1 (47.118.XXX.XXX/47.118.XXX.XXX:18001) could not be established. Broker may not be available.
    [2025-02-25 18:22:27.618] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ClusterSwitchListener:44] - Cluster not changed on update:5aPLLlDtTHqP8sKq-DZVfg
    ******

    SDK クライアントは定期的にデータ消費の統計情報を収集して表示します。これらの統計情報には、送受信されたデータレコードの総数、データの総量、および 1 秒あたりのレコード数 (RPS) が含まれます。

    [2025-02-25 18:22:18.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174047893****,"DefaultUserRecordQueue":0.0}

    パラメータ

    説明

    outCounts

    SDK クライアントが消費したデータレコードの総数です。

    outBytes

    SDK クライアントが消費したデータの総量 (バイト単位) です。

    outRps

    SDK クライアントによるデータ消費の 1 秒あたりのレコード数です。

    outBps

    SDK クライアントによるデータ消費の 1 秒あたりに転送されるバイト数です。

    count

    データ消費情報 (メトリクス) 内のパラメータの総数です。

    説明

    これには count 自体は含まれません。

    inBytes

    DTS サーバーが送信したデータの総量 (バイト単位) です。

    DStoreRecordQueue

    DTS サーバーがデータを送信する際のデータキャッシュキューの現在のサイズです。

    inCounts

    DTS サーバーが送信したデータレコードの総数です。

    inBps

    DTS サーバーがデータを送信する際の 1 秒あたりに転送されるバイト数です。

    inRps

    DTS サーバーがデータを送信する際の 1 秒あたりのレコード数です。

    __dt

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

    DefaultUserRecordQueue

    シリアル化後のデータキャッシュキューのサイズです。

  6. 必要に応じて、サブスクライブしたデータを消費するためのコードを編集します。

    サブスクライブしたデータを消費する際は、データ損失を防ぎ、データの重複を最小限に抑え、オンデマンドの消費を可能にするために、コンシューマーオフセットを管理する必要があります。

よくある質問

  • サブスクリプションインスタンスに接続できない場合はどうすればよいですか?

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

  • 永続化後のコンシューマーオフセットのデータ形式は何ですか?

    コンシューマーオフセットが永続化されると、データは JSON 形式で返されます。永続化されたコンシューマーオフセットは、SDK に直接渡すことができる Unix タイムスタンプです。たとえば、返されたデータでは、"timestamp" キーの値 1700709977 が永続化されたコンシューマーオフセットです。

    {"groupID":"dtsglg11d48230***","streamCheckpoint":[{"partition":0,"offset":577989,"topic":"ap_southeast_1_vpc_rm_t4n22s21iysr6****_root_version2","timestamp":1700709977,"info":""}]}
  • 変更追跡タスクを複数のクライアントで並列に消費できますか?

    いいえ。SUBSCRIBE モードでは複数のクライアントを並列に実行できますが、一度に 1 つのクライアントのみがデータを消費できます。

  • SDK コードにカプセル化されている Kafka クライアントのバージョンは何ですか?

    dts-new-subscribe-sdk のバージョン 2.0.0 以降は、Kafka クライアント (kafka-clients) 2.7.0 をカプセル化しています。2.0.0 より前のバージョンは Kafka クライアント 1.0.0 をカプセル化しています。

    説明

    アプリケーション開発プロセスで依存関係パッケージの脆弱性検出ツールを使用し、dts-new-subscribe-sdk にカプセル化されている Kafka クライアント (kafka-clients) にセキュリティ脆弱性があることが判明した場合、クライアントを 2.1.4-shaded バージョンに置き換えることで、この脆弱性を解決できます。

    <dependency>
        <groupId>com.aliyun.dts</groupId>
        <artifactId>dts-new-subscribe-sdk</artifactId>
        <version>2.1.4-shaded</version>
    </dependency>

付録

コンシューマーオフセットの管理

SDK クライアントが初めて起動される場合、再起動される場合、または内部的に再試行される場合、データ消費を開始または再開するには、コンシューマーオフセットをクエリして渡す必要があります。コンシューマーオフセットは、SDK クライアントが消費する最初のデータレコードの UNIX タイムスタンプです。

クライアントのコンシューマーオフセットをリセットするには、次の表で説明されている消費モード (SDK 使用モード) に基づいて、コンシューマーオフセットをクエリおよび変更できます。

説明

上記にリストされているチェックポイントストレージレベル (外部ストレージ、localCheckpointStore ファイル、DTS Server/DStore オフセット) は、すべて単一のコミット操作で書き込まれるため、その内容はほぼ同一です。主な違いは、localCheckpointStore ファイルはノードまたはコンテナが破棄されると失われるのに対し、DTS Server (DStore) オフセットはサーバー側で永続化される点です。

シナリオ

SDK 使用モード

オフセット管理方法

コンシューマーオフセットのクエリ

ASSIGN モード、SUBSCRIBE モード

  • SDK クライアントはメッセージオフセットを 5 秒ごとに保存し、DTS サーバーにコミットするため、最後のコンシューマーオフセットを次のいずれかの場所からクエリできます:

    • SDK クライアントが稼働しているサーバー上の localCheckpointStore ファイル。

    • 変更追跡チャネルの[データ消費] ページ

  • consumerContext.java ファイルで setUserRegisteredStore(new UserMetaStore()) を使用してデータベースなどの外部永続化共有ストレージメディアを構成している場合、このストレージメディアは 5 秒ごとにメッセージオフセットを保存し、クエリできます。

SDK クライアントが初めて起動される場合。データを消費するためにコンシューマーオフセットを渡す必要があります。

ASSIGN モード、SUBSCRIBE モード

SDK クライアントの使用パターンに応じて、Java ファイル DTSConsumerAssignDemo.java または DTSConsumerSubscribeDemo.java を選択し、コンシューマーオフセットを構成 (initCheckpoint) してデータを消費します。

SDK クライアントが内部的に再試行される場合。データ消費を再開するために、最後に記録されたコンシューマーオフセットを渡す必要があります。

ASSIGN モード

次の順序で最後に記録されたコンシューマーオフセットを検索します。オフセットが見つかった場合、オフセット情報が返されます。

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

  2. SDK クライアントが稼働しているサーバー上の localCheckpointStore ファイル。

  3. DTS Server (DStore) に保存されているオフセット。

  4. DTSConsumerAssignDemo.java ファイルの initCheckpoint に渡した開始タイムスタンプ。

SUBSCRIBE モード

次の順序で最後に記録されたコンシューマーオフセットを検索します。オフセットが見つかった場合、オフセット情報が返されます。

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

  2. DTS Server (増分データインジェストモジュール) に保存されているオフセット。

    説明

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

  3. DTSConsumerSubscribeDemo.java ファイルの initCheckpoint に渡した開始タイムスタンプ。

  4. DTS Server (新しい増分データインジェストモジュール) の開始オフセット。

    重要

    増分データインジェストモジュールが切り替わった場合、新しいモジュールはクライアントの最後のコンシューマーオフセットを保存できません。これにより、データ消費が古いオフセットから開始される可能性があります。クライアント側でコンシューマーオフセットを永続化して保存することを推奨します。

SDK クライアントが再起動される場合。データ消費を再開するために、最後に記録されたコンシューマーオフセットを渡す必要があります。

ASSIGN モード

consumerContext.java ファイルの setForceUseCheckpoint 構成に基づいて、コンシューマーオフセットがクエリされ、見つかった場合はオフセット情報が返されます。

  • true に設定されている場合、SDK クライアントは再起動されるたびに渡された initCheckpoint をコンシューマーオフセットとして使用します。

  • false として構成されているか、構成されていない場合、次の順序で前のレコードのコンシューマーオフセットを検索します。

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

    2. SDK クライアントが稼働しているサーバー上の localCheckpointStore ファイル。

    3. DTS Server (増分データインジェストモジュール) に保存されているオフセット。

      説明

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

SUBSCRIBE モード

このモードでは、consumerContext.java ファイルの setForceUseCheckpoint 構成は有効になりません。次の順序で前のレコードのコンシューマーオフセットを検索します。

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

  2. DTS Server (増分データインジェストモジュール) に保存されているオフセット。

    説明

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

  3. DTSConsumerSubscribeDemo.java ファイルの initCheckpoint に渡した開始タイムスタンプ。

  4. DTS Server (新しい増分データインジェストモジュール) の開始オフセット。

コンシューマーオフセットの永続化

増分データインジェストモジュールのディザスタリカバリ切り替えが発生した場合、新しいモジュールはクライアントの最後のコンシューマーオフセットを保存できません。これは特に 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. AbstractUserMetaStore() メソッドを継承して実装する 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();
                              
                if (rs.next()) {
                    String checkpoint = rs.getString("checkpoint");
                    return checkpoint;
                }
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                close(rs, pres, con);
            }
            return null;
        }
    }
    
  2. consumerContext.java ファイルで、setUserRegisteredStore(new UserMetaStore()) メソッドを呼び出して外部ストレージメディアを構成します。

トラブルシューティング

例外

エラーメッセージ

原因

ソリューション

接続失敗

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 ファイルで setForceUseCheckpoint が true に設定されていますが、コンシューマーオフセットがサブスクリプションインスタンスのタイムスタンプ範囲内にありません。

サブスクリプションインスタンスのタイムスタンプ範囲内のコンシューマーオフセットを入力してください。クエリ方法については、「パラメータの説明」をご参照ください。

データ消費の低速化

N/A

  • DStoreRecordQueue と DefaultUserRecordQueue のキューのサイズに関する統計情報の各パラメータをクエリし、データ消費が低速化した理由を分析します。

    • DStoreRecordQueue の値が 0 のままの場合、DTS サーバーがデータをプルする速度が遅くなっています。

    • DefaultUserRecordQueue の値がデフォルト値の 512 のままの場合、SDK クライアントがデータを消費する速度が遅くなっています。

  • 必要に応じて、コード内のコンシューマーオフセット (initCheckpoint) を変更してオフセットをリセットします。