変更追跡タスクを作成した後、Data Transmission Service (DTS) が提供する SDK (Software Development Kit) を使用して、データ変更をサブスクライブできます。このトピックでは、SDK を使用して、PolarDB-X 1.0 や DMS 論理データベースなどの分散データソースからデータを使用する方法について説明します。
前提条件
変更追跡インスタンスが作成されました。 インスタンスは 通常 状態です。 詳細については、「PolarDB-X 1.0 インスタンスの変更追跡タスクを作成する」または「DMS 論理データベースの変更追跡タスクを作成する」をご参照ください。
サブスクリプションインスタンスの使用者グループを作成しました。
RAM ユーザーを使用してサブスクライブデータをコンシュームする場合、その RAM ユーザーには AliyunDTSFullAccess 権限と、サブスクライブ対象オブジェクトへのアクセス権限が必要です。詳細については、「システムポリシーを使用して DTS の管理権限を RAM ユーザーに付与する」および「RAM ユーザーの権限管理」をご参照ください。
注意事項
サブスクライブデータをコンシュームする際は、`DefaultUserRecord` の commit メソッドを呼び出してオフセット情報をコミットする必要があります。そうしないと、データが繰り返しコンシュームされる可能性があります。
異なるコンシュームプロセスは互いに独立しています。
手順
SDK サンプルコード をダウンロードして解凍します。
SDK コードのバージョンを確認します。
サンプル SDK コードを解凍したディレクトリに移動します。
テキストエディターで、ディレクトリ内の pom.xml ファイルを開きます。
変更追跡 SDK を最新バージョンに更新します。
説明最新の Maven 依存関係は、dts-new-subscribe-sdk ページで確認できます。
SDK コードを編集します。
解凍したファイルを IDE またはテキストエディタで開くことができます。
SDK クライアントの使用パターンに基づいて、DistributedDTSConsumerDemo.java ファイルを開きます。
説明Java ファイルのパスは
aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/です。Java コード内のパラメータを設定します。
public static void main(String[] args) throws ClientException { // PolarDB-X 1.0 (旧 DRDS) などの分散データソースをサブスクライブするための設定。AccessKey、インスタンスID、タスクID、コンシューマーグループなどの情報を設定します。 String accessKeyId = "LTA***********99reZ"; String accessKeySecret = "****************"; String regionId = "cn-hangzhou"; String dtsInstanceId = "dtse5212sed162****"; String jobId = "l791216x16d****"; String sid = "dtsip412t13160****"; String userName = "xftest"; String password = "******"; String proxyUrl = "dts-cn-****.com:18001"; // 初回シーク用の初期チェックポイント (設定するタイムスタンプ、例:2019 年 8 月 19 日 10:03:21 CST が必要な場合は 1566180200) String checkpoint = "1639620090"; // 物理データベース/テーブル名を論理データベース/テーブル名に変換 boolean mapping = true; // 開始時に設定チェックポイントを強制的に使用するかどうか。チェックポイントのリセットの場合、割り当てモードのみが機能します boolean isForceUseInitCheckpoint = false; ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.ASSIGN; DistributedDTSConsumerDemo demo = new DistributedDTSConsumerDemo(userName, password, regionId, jobId, sid, dtsInstanceId, accessKeyId, accessKeySecret, subscribeMode, proxyUrl, checkpoint, isForceUseInitCheckpoint, mapping); demo.start(); }パラメータ
説明
取得方法
accessKeyId
AccessKey ID。
詳細については、「AccessKey ペアを取得する」をご参照ください。
accessKeySecret
AccessKey シークレット。
regionId
変更追跡タスクが配置されているリージョンの ID。
DTS コンソールで、ターゲットの変更追跡インスタンスの ID をクリックします。基本情報 ページで、リージョン情報を取得できます。たとえば、リージョンが中国 (杭州) の場合、このパラメーターを
cn-hangzhouに設定します。詳細については、「リージョンの一覧」をご参照ください。dtsInstanceId
変更追跡インスタンスの ID。
DTS コンソールで、ターゲットの変更追跡インスタンスの ID をクリックします。基本情報 ページで、変更追跡インスタンスの DTS インスタンス ID を取得できます。
jobId
変更追跡タスクの ID。
DescribeDtsJobs 操作を呼び出して、変更追跡タスク ID ([DtsJobId]) を取得できます。
sid
コンシューマーグループの ID。
DTS コンソールで、対象の変更追跡インスタンスの ID をクリックします。左側のナビゲーションペインで、データ消費をクリックします。コンシューマーグループのコンシューマーグループ ID /名前とアカウントを取得できます。
説明コンシューマーグループアカウントのパスワードは、コンシューマーグループの作成時に指定されます。
userName
コンシューマーグループのアカウント。
password
コンシューマーグループアカウントのパスワード。
proxyUrl
変更追跡チャネルのエンドポイントとポート。
説明SDK クライアントをデプロイする ECS インスタンスと変更追跡チャネルが同じクラシックネットワークまたは仮想プライベートクラウド (VPC) 内にある場合、内部ネットワーク経由でデータをサブスクライブすることで、最小のレイテンシを実現できます。
ネットワークが不安定になる可能性があるため、パブリックエンドポイントの使用は推奨されません。
DTS コンソールで、対象の変更追跡インスタンスの ID をクリックします。基本情報 ページでは、ネットワーク 情報を取得できます。
checkpoint
コンシューマーオフセット。SDK クライアントがデータレコードの使用を開始するタイムスタンプです。値は UNIX タイムスタンプ (秒単位) です。
説明コンシューマーオフセット情報は以下の用途に使用できます:
使用プロセスが中断された場合、コンシューマーオフセットを渡してデータの使用を再開し、データ損失を防ぐことができます。
SDK クライアントを起動する際、必要なコンシューマーオフセットを渡してサブスクリプションオフセットを調整し、必要に応じてデータを使用できます。
コンシューマーオフセットは、変更追跡インスタンスのタイムスタンプ範囲内である必要があり、UNIX タイムスタンプに変換する必要があります。
説明追跡タスクリストのデータ範囲列で、変更追跡インスタンスのタイムスタンプ範囲を表示できます。
UNIX タイムスタンプコンバーターは、検索エンジンなどで見つけることができます。
オプション:サブスクライブしたデータのデータ型を変更するには、
buildRecordListener()メソッドを変更するか、カスタムクラスを使用します。public static Map<String, RecordListener> buildRecordListener() { // ユーザーは独自のリスナーを実装できます RecordListener mysqlRecordPrintListener = new RecordListener() { @Override public void consume(DefaultUserRecord record) { OperationType operationType = record.getOperationType(); if (operationType.equals(OperationType.INSERT) || operationType.equals(OperationType.UPDATE) || operationType.equals(OperationType.DELETE) || operationType.equals(OperationType.HEARTBEAT)) { // レコードを使用 RecordListener recordPrintListener = new DefaultRecordPrintListener(DbType.MySQL); recordPrintListener.consume(record); // commit メソッドでチェックポイントを更新します record.commit(""); } } }; return Collections.singletonMap("mysqlRecordPrinter", mysqlRecordPrintListener); }IDE でプロジェクト構造を開き、プロジェクトの OpenJDK バージョンが 1.8 であることを確認します。
クライアントコードを実行します。
出力には、クライアントがソースデータベースからのデータ変更をサブスクライブしていることが表示されます。
SDK クライアントは、データ使用に関する統計情報を定期的に収集して表示します。統計情報には、送受信されたデータレコードの総数、データの総量、および 1 秒あたりのリクエスト数 (RPS) が含まれます。
表 1. データ使用統計
パラメータ
説明
outCountsSDK クライアントが使用したデータレコードの総数。
outBytesSDK クライアントが使用したデータの総量 (バイト単位)。
outRpsSDK クライアントがデータを使用するために送信する 1 秒あたりのリクエスト数。
outBpsSDK クライアントがデータを使用する際に 1 秒あたりに転送されるビット数。
countなし。
inBytesDTS サーバーが送信したデータの総量 (バイト単位)。
DStoreRecordQueueDTS サーバーがデータを送信する際のデータキャッシュキューのサイズ。
inCountsDTS サーバーが送信したデータレコードの総数。
inRpsDTS サーバーが 1 秒あたりに送信するリクエストの数。
inBpsDTS サーバーがデータを送信する際に 1 秒あたりに転送されるビット数。
__dtSDK クライアントがデータを受信した際のタイムスタンプ (ミリ秒単位)。
DefaultUserRecordQueueシリアル化後のデータキャッシュキューのサイズ。