データサブスクリプションインスタンスを設定した後、 Data Transmission Service (DTS) の SDK サンプルコードを使用して変更データを消費します。
操作手順
-
データソースが PolarDB-X 1.0 インスタンスまたは DMS 論理データベースの場合、「SDK サンプルコードを使用した PolarDB-X 1.0 からのデータサブスクリプションデータの消費」をご参照ください。
-
RAM ユーザーを使用してデータを消費する場合、RAM ユーザーには AliyunDTSFullAccess 権限と、サブスクリプションオブジェクトにアクセスする権限が必要です。権限の付与方法の詳細については、「システムポリシーを使用した DTS インスタンスを管理するための RAM ユーザーの承認」および「RAM ユーザーの権限の管理」をご参照ください。
-
各コンシューマーは独立して動作します。
-
このトピックでは、Java のサンプル SDK クライアントを提供します。Python と Go のサンプルコードについては、「dts-subscribe-demo」をご参照ください。
次の手順では、IntelliJ IDEA (Community Edition 2020.1 for Windows) で SDK サンプルコードを実行して、データサブスクリプションのデータを消費する方法について説明します。
-
データサブスクリプションインスタンスを作成します。詳細については、「ApsaraDB RDS for MySQL インスタンスのデータサブスクリプションチャネルの作成」、「PolarDB for MySQL クラスターのデータサブスクリプションチャネルの作成」、または「Oracle データベースのデータサブスクリプションチャネルの作成」をご参照ください。
-
1 つ以上のコンシューマーグループを作成します。詳細については、「コンシューマーグループの作成」をご参照ください。
重要データサブスクリプションのデータを消費する場合、
DefaultUserRecordのcommitメソッドを呼び出してチェックポイントをコミットする必要があります。そうしないと、データが重複して消費される可能性があります。 -
ビジネス要件に基づいて SDK サンプルコードを使用します。
-
新しいデータサブスクリプション SDK パッケージを使用する (推奨)
-
IntelliJ IDEA を開き、[Create New Project] をクリックして、アプリケーション用のプロジェクトを作成します。
-
プロジェクトで、プロジェクトオブジェクトモデル (POM) ファイル pom.xml を見つけます。
-
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 のページで確認できます。
-
新しいサブスクリプション SDK の使用方法の詳細については、「サンプルコードの使用」をご参照ください。
-
-
新しいデータサブスクリプション SDK のカスタマイズ版を使用する
-
SDK サンプルコードパッケージをダウンロードし、解凍します。
説明
をクリックし、[Download ZIP] を選択してパッケージをダウンロードします。 -
解凍した 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 ページ」をご参照ください。
-
IntelliJ IDEA を開きます。ようこそ画面で、[Open or Import] をクリックします。
-
[Open File or Project] ダイアログボックスで、解凍した SDK サンプルコードのディレクトリに移動し、pom.xml ファイルを選択して [OK] をクリックします。
-
表示されるダイアログボックスで、[Open as Project] を選択します。
-
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 に割り当てられ、消費を再開します。
-
-
-
-
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ユーザー名のパスワード。
initCheckpointSDK クライアントがデータの消費を開始する消費チェックポイント。UNIX タイムスタンプとして指定します。例:1620962769。
説明次のシナリオで消費チェックポイント情報を使用できます:
-
アプリケーションが中断した後に消費を再開し、データの損失を防ぐには、最後に認識された消費チェックポイントを渡してください。
-
クライアントを起動するときに、特定の消費チェックポイントを渡して、目的の位置からデータを消費できます。
消費チェックポイントは、データサブスクリプションインスタンスのデータ範囲 (図を参照) 内にある必要があり、UNIX タイムスタンプに変換する必要があります。
説明検索エンジンを使用して、UNIX タイムスタンプコンバーターを見つけることができます。
ConsumerContext.ConsumerSubscribeMode subscribeModeSDK クライアントの使用モード。有効な値:
-
ConsumerContext.ConsumerSubscribeMode.ASSIGN:ASSIGN モード。コンシューマーグループ内の 1 つの SDK クライアントのみがデータサブスクリプションのデータを消費できます。 -
ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE:SUBSCRIBE モード。災害復旧のために、同じコンシューマーグループで複数の SDK クライアントを起動できます。
N/A
-
-
IntelliJ IDEA の上部メニューで、 を選択してクライアントを実行します。
説明クライアントを初めて実行するときは、必要な依存関係の自動ロードとインストールに時間がかかる場合があります。
-
実行後、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. データ消費統計
パラメーター
説明
outCountsSDK クライアントによって消費されたデータレコードの総数。
outBytesSDK クライアントによって消費されたデータの総量。単位:バイト。
outRpsSDK クライアントがデータを消費する際の RPS (Requests Per Second)。
outBpsSDK クライアントのデータ消費レート。単位:バイト/秒。
inBytesDTS サーバーから送信されたデータの総量。単位:バイト。
DStoreRecordQueueDTS サーバーから受信したレコードを格納する内部データキャッシュキューのサイズ。
inCountsDTS サーバーから送信されたデータレコードの総数。
inRpsDTS サーバーがデータを送信する際の RPS。
__dtSDK クライアントがデータを受信したときのタイムスタンプ。単位:ミリ秒。
DefaultUserRecordQueueコンシューマーアプリケーションが処理する準備ができたレコードを保持するデータキューのサイズ。
-
消費チェックポイントの保存とクエリ
データ消費を開始または再開する際 (初回起動時、再起動時、内部リトライ時など)、SDK クライアントには消費チェックポイントが必要です。次の表では、データ損失を防ぎ、重複消費を最小限に抑え、オンデマンド消費を可能にするために、さまざまなシナリオでチェックポイントを管理およびクエリする方法について説明します。
|
シナリオ |
SDK 利用モード |
クエリ方法 |
|
消費チェックポイントのクエリ |
ASSIGN モード、SUBSCRIBE モード |
|
|
初回起動: チェックポイントを渡して消費を開始する。 |
ASSIGN モード、SUBSCRIBE モード |
SDK クライアントの利用モードに基づいて、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) から開始するため、重複消費が発生します。
このシナリオで重複消費を回避するには、クライアントに永続的なチェックポイントストアを設定してください。次の例では、要件に応じて調整できる実装例を 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); } } } -
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":""}]}
トラブルシューティング
|
問題 |
エラーメッセージ |
原因 |
解決策 |
|
接続できません |
|
指定された |
|
|
ブローカーアドレスが実際の IP アドレスに接続できません。 |
||
|
ユーザー名またはパスワードが正しくありません。 |
||
|
consumerContext.java ファイルで、 |
データサブスクリプションインスタンスのデータ範囲内にある消費チェックポイントを渡します。詳細については、「必須パラメーター」をご参照ください。 |
|
|
消費速度の低下 |
N/A |
|
|