変更追跡タスクとコンシューマーグループを作成して変更追跡チャネルを構成した後、Data Transmission Service (DTS) が提供するソフトウェア開発キット (SDK) を使用して、サブスクライブしたデータを消費できます。このトピックでは、サンプルコードの使用方法について説明します。
PolarDB-X 1.0 データソースからサブスクライブしたデータを消費する場合は、「SDK を使用した PolarDB-X 1.0 からのサブスクライブ済みデータの消費」をご参照ください。
このトピックでは、Java の SDK クライアント向けのサンプルコードを提供しています。Python および Go のサンプルコードについては、「dts-subscribe-demo」をご参照ください。
前提条件
サブスクリプションインスタンスが作成され、通常 状態で稼働しています。
説明サブスクリプションインスタンスの作成手順については、「サブスクリプションプランの概要」をご参照ください。
サブスクリプションインスタンス用のコンシューマーグループを作成していること。
RAM ユーザーを使用してサブスクライブしたデータを消費する場合、その RAM ユーザーには AliyunDTSFullAccess 権限と、サブスクライブ対象のオブジェクトへのアクセス権限が必要です。詳細については、「システムポリシーを使用した RAM ユーザーへの DTS 管理権限の付与」および「RAM ユーザーの権限管理」をご参照ください。
注意事項
サブスクライブしたデータを消費する際は、データ消費を完了してから
DefaultUserRecordのcommitメソッドを呼び出してオフセット情報をコミットする必要があります。データを消費する前に commit メソッドを呼び出さないでください。そうしないと、データ損失が発生します。異なる消費プロセスは互いに独立しています。
コンソールでは、[現在のオフセット] は、追跡タスクがサブスクライブしたオフセットを示しており、クライアントによってコミットされたオフセットではありません。
操作手順
サンプル SDK コードファイルをダウンロードし、パッケージを解凍します。
SDK コードのバージョンを確認します。
サンプル SDK コードを解凍したディレクトリに移動します。
テキストエディタを使用して、ディレクトリ内の pom.xml ファイルを開きます。
変更追跡 SDK を最新バージョンに更新します。
説明最新の Maven 依存関係は、「dts-new-subscribe-sdk」ページで確認できます。
SDK コードを編集します。
統合開発環境 (IDE) を使用して、解凍したファイルを開きます。
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 クライアントがサブスクライブしたデータを消費します。これはディザスタリカバリシナリオです。
Java コード内のパラメータを設定します。
パラメータ
説明
取得方法
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 タイムスタンプに変換する必要があります。
説明追跡タスク一覧にあるデータ範囲列には、ターゲットサブスクリプションインスタンスのタイムスタンプの範囲が表示されます。
subscribeModeSDK クライアントの使用モードです。このパラメータを変更する必要はありません。
ConsumerContext.ConsumerSubscribeMode.ASSIGN: ASSIGN モードConsumerContext.ConsumerSubscribeMode.SUBSCRIBE: SUBSCRIBE モード
N/A
IDE でプロジェクト構造を開き、プロジェクトの OpenJDK バージョンが 1.8 であることを確認します。
クライアントコードを実行します。
説明初回実行時、IDE は必要なプラグインと依存関係を自動的にロードするため、時間がかかることがあります。
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}パラメータ
説明
outCountsSDK クライアントが消費したデータレコードの総数です。
outBytesSDK クライアントが消費したデータの総量 (バイト単位) です。
outRpsSDK クライアントによるデータ消費の 1 秒あたりのレコード数です。
outBpsSDK クライアントによるデータ消費の 1 秒あたりに転送されるバイト数です。
countデータ消費情報 (メトリクス) 内のパラメータの総数です。
説明これには
count自体は含まれません。inBytesDTS サーバーが送信したデータの総量 (バイト単位) です。
DStoreRecordQueueDTS サーバーがデータを送信する際のデータキャッシュキューの現在のサイズです。
inCountsDTS サーバーが送信したデータレコードの総数です。
inBpsDTS サーバーがデータを送信する際の 1 秒あたりに転送されるバイト数です。
inRpsDTS サーバーがデータを送信する際の 1 秒あたりのレコード数です。
__dtSDK クライアントがデータを受信したときのタイムスタンプ (ミリ秒単位) です。
DefaultUserRecordQueueシリアル化後のデータキャッシュキューのサイズです。
必要に応じて、サブスクライブしたデータを消費するためのコードを編集します。
サブスクライブしたデータを消費する際は、データ損失を防ぎ、データの重複を最小限に抑え、オンデマンドの消費を可能にするために、コンシューマーオフセットを管理する必要があります。
よくある質問
サブスクリプションインスタンスに接続できない場合はどうすればよいですか?
エラーメッセージに基づいて問題をトラブルシューティングしてください。詳細については、「トラブルシューティング」をご参照ください。
永続化後のコンシューマーオフセットのデータ形式は何ですか?
コンシューマーオフセットが永続化されると、データは 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 クライアントが初めて起動される場合。データを消費するためにコンシューマーオフセットを渡す必要があります。 | ASSIGN モード、SUBSCRIBE モード | SDK クライアントの使用パターンに応じて、Java ファイル DTSConsumerAssignDemo.java または DTSConsumerSubscribeDemo.java を選択し、コンシューマーオフセットを構成 ( |
SDK クライアントが内部的に再試行される場合。データ消費を再開するために、最後に記録されたコンシューマーオフセットを渡す必要があります。 | ASSIGN モード | 次の順序で最後に記録されたコンシューマーオフセットを検索します。オフセットが見つかった場合、オフセット情報が返されます。
|
SUBSCRIBE モード | 次の順序で最後に記録されたコンシューマーオフセットを検索します。オフセットが見つかった場合、オフセット情報が返されます。
| |
SDK クライアントが再起動される場合。データ消費を再開するために、最後に記録されたコンシューマーオフセットを渡す必要があります。 | ASSIGN モード | consumerContext.java ファイルの
|
SUBSCRIBE モード | このモードでは、consumerContext.java ファイルの
|
コンシューマーオフセットの永続化
増分データインジェストモジュールのディザスタリカバリ切り替えが発生した場合、新しいモジュールはクライアントの最後のコンシューマーオフセットを保存できません。これは特に 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) から消費を開始し、データの重複消費が発生します。
この切り替えシナリオで履歴データの重複消費を防ぐため、クライアント側でコンシューマーオフセットの永続化ストレージを構成することを推奨します。以下に参考用のサンプルメソッドを示します。必要に応じて変更できます。
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; } }consumerContext.java ファイルで、
setUserRegisteredStore(new UserMetaStore())メソッドを呼び出して外部ストレージメディアを構成します。
トラブルシューティング
例外 | エラーメッセージ | 原因 | ソリューション |
接続失敗 | |
| 正しい |
| ブローカーアドレスを使用して実際の IP アドレスに接続できません。 | ||
| ユーザー名またはパスワードが正しくありません。 | ||
| consumerContext.java ファイルで | サブスクリプションインスタンスのタイムスタンプ範囲内のコンシューマーオフセットを入力してください。クエリ方法については、「パラメータの説明」をご参照ください。 | |
データ消費の低速化 | N/A |
| |