Data Transmission Service (DTS) は、ApsaraDB for MongoDB のシャーディングクラスターから増分変更をキャプチャし、各変更イベントを Canal JSON 形式で Function Compute の関数に配信します。関数コードを記述して、イベントを処理、変換、またはダウンストリームに転送します。
関数が受信する内容
各呼び出しは、Records 配列を持つオブジェクトを渡します。配列内のすべての要素は、1 つの変更イベントを表します。
| フィールド | タイプ | 説明 |
|---|---|---|
isDdl |
Boolean | DDL 操作の場合は True、DML 操作の場合は False |
type |
String | DML: INSERT、UPDATE、または DELETE。DDL: DDL |
database |
String | MongoDB データベース名 |
table |
String | コレクション名 |
pkNames |
String | プライマリキー名。MongoDB の場合は常に _id です |
es |
Long | 13 桁の UNIX タイムスタンプ (ミリ秒) — ソースで操作が発生した時刻 |
ts |
Long | 13 桁の UNIX タイムスタンプ (ミリ秒) — DTS が宛先への書き込みを開始した時刻 |
data |
Object Array | 要素が 1 つの配列。要素には doc キーがあり、その値はドキュメントを表す JSON 文字列です。値をデシリアライズしてレコードを読み取ります |
old |
Object Array | data と同じ形式です。type が UPDATE の場合にのみ存在し、更新前のドキュメントの状態が含まれます |
id |
Int | 操作のシリアル番号 |
関数は 2 つのカテゴリの操作を受信します。
-
DDL — スキーマ変更:
CreateIndex、CreateCollection、DropIndex、DropCollection -
DML — データ変更:
INSERT、UPDATE、DELETE
ペイロードの例
コレクションの作成 (DDL)
db.createCollection("testCollection"){
"Records": [{
"data": [{"doc": "{\"create\": \"testCollection\", \"idIndex\": {\"v\": 2, \"key\": {\"_id\": 1}, \"name\": \"_id_\"}}"}],
"pkNames": ["_id"],
"type": "DDL",
"es": 1694056437000,
"database": "MongoDBTest",
"id": 0,
"isDdl": true,
"table": "testCollection",
"ts": 1694056437510
}]
}
コレクションの削除 (DDL)
db.testCollection.drop(){
"Records": [{
"data": [{"doc": "{\"drop\": \"testCollection\"}"}],
"pkNames": ["_id"],
"type": "DDL",
"es": 1694056577000,
"database": "MongoDBTest",
"id": 0,
"isDdl": true,
"table": "testCollection",
"ts": 1694056577789
}]
}
インデックスの作成 (DDL)
db.testCollection.createIndex({name: 1}){
"Records": [{
"data": [{"doc": "{\"createIndexes\": \"testCollection\", \"v\": 2, \"key\": {\"name\": 1}, \"name\": \"name_1\"}"}],
"pkNames": ["_id"],
"type": "DDL",
"es": 1694056670000,
"database": "MongoDBTest",
"id": 0,
"isDdl": true,
"table": "testCollection",
"ts": 1694056670719
}]
}
インデックスの削除 (DDL)
db.testCollection.dropIndex({name: 1}){
"Records": [{
"data": [{"doc": "{\"dropIndexes\": \"testCollection\", \"index\": \"name_1\"}"}],
"pkNames": ["_id"],
"type": "DDL",
"es": 1694056817000,
"database": "MongoDBTest",
"id": 0,
"isDdl": true,
"table": "$cmd",
"ts": 1694056818035
}]
}
ドキュメントの挿入 (DML)
// バッチ挿入
db.runCommand({insert: "user", documents: [{"name": "jack", "age": 20}, {"name": "lili", "age": 20}]})
// 単一挿入
db.user.insert({"name": "jack", "age": 20}){
"Records": [
{
"data": [{"doc": "{\"_id\": {\"$oid\": \"64f9397f6e255f74d65a****\"}, \"name\": \"jack\", \"age\": 20}"}],
"pkNames": ["_id"],
"type": "INSERT",
"es": 1694054783000,
"database": "MongoDBTest",
"id": 0,
"isDdl": false,
"table": "user",
"ts": 1694054784427
},
{
"data": [{"doc": "{\"_id\": {\"$oid\": \"64f9397f6e255f74d65a****\"}, \"name\": \"lili\", \"age\": 20}"}],
"pkNames": ["_id"],
"type": "INSERT",
"es": 1694054783000,
"database": "MongoDBTest",
"id": 0,
"isDdl": false,
"table": "user",
"ts": 1694054784428
}
]
}
ドキュメントの更新 (DML)
db.user.update({"name": "jack"}, {$set: {"age": 30}}){
"Records": [{
"data": [{"doc": "{\"$set\": {\"age\": 30}}"}],
"pkNames": ["_id"],
"old": [{"doc": "{\"_id\": {\"$oid\": \"64f9397f6e255f74d65a****\"}}"}],
"type": "UPDATE",
"es": 1694054989000,
"database": "MongoDBTest",
"id": 0,
"isDdl": false,
"table": "user",
"ts": 1694054990555
}]
}
UPDATE 操作では、DTS が増分データを同期するときに $set コマンドのみが同期的に実行されます。
ドキュメントの削除 (DML)
db.user.remove({"name": "jack"}){
"Records": [{
"data": [{"doc": "{\"_id\": {\"$oid\": \"64f9397f6e255f74d65a****\"}}"}],
"pkNames": ["_id"],
"type": "DELETE",
"es": 1694055452000,
"database": "MongoDBTest",
"id": 0,
"isDdl": false,
"table": "user",
"ts": 1694055452852
}]
}
制限事項
同期タスクを作成する前に、これらの制約を確認してください。
スコープの制約:
-
増分データ同期のみがサポートされます。フルデータ同期はサポートされません。
-
クロスリージョン同期はサポートされません。
-
DTS は、
admin、config、またはlocalデータベースからデータを同期できません。 -
オブジェクトマッピングはサポートされません。
-
トランザクション情報は保持されません。トランザクションは、宛先で個別のレコードに変換されます。
ソースデータベースの制約:
-
ソースはシャーディングクラスターアーキテクチャを使用する必要があり、Mongos ノードは 10 個以下である必要があります。
-
ソースを Azure Cosmos DB for MongoDB クラスターまたは Amazon DocumentDB エラスティッククラスターにすることはできません。
-
同期するコレクションには、プライマリキーまたは UNIQUE 制約が必要で、すべてのフィールドが一意である必要があります。そうでない場合、宛先データベースに重複したデータレコードが含まれる可能性があります。
-
単一のドキュメントは 16 MB を超えることはできません。このサイズを超えるドキュメントは宛先の関数に書き込むことができず、エラーを引き起こします。必要に応じて、抽出、変換、ロード (ETL) 機能を使用して大きなフィールドをフィルターします。
-
タスクごとに最大 1,000 個のコレクションを同期できます。より多くのコレクションを同期するには、複数のタスクを作成するか、データベースレベルで同期します。
-
同期タスクの実行中は、ソースインスタンスをスケーリングできません。
-
DTS は SRV エンドポイント経由で MongoDB データベースに接続できません。
-
ソースデータベースのバランサーがアクティブな場合、タスクに遅延が発生することがあります。
-
ソースデータベースがシャーディングクラスターアーキテクチャを使用するセルフマネージド MongoDB データベースの場合、[アクセス方法] パラメーターを Express Connect、VPN Gateway、または Smart Access Gateway または Cloud Enterprise Network (CEN) に設定します。
書き込みの制約:
-
INSERT 操作の場合、挿入されるデータにはシャードキーが含まれている必要があります。
-
UPDATE 操作の場合、シャードキーは変更できません。
-
宛先の関数ごとに 1 つの DTS タスクのみを設定してください。同じ関数に書き込む複数のタスクは、データエラーを引き起こす可能性があります。
oplogと変更ストリームの要件:
-
oplog を有効にして少なくとも 7 日間のログデータを保持するか、変更ストリームを有効にして少なくとも過去 7 日間の変更をカバーする必要があります。どちらの条件も満たされない場合、DTS はソースの変更をキャプチャできず、DTS のサービスレベルアグリーメント (SLA) の対象外となるデータの不整合や損失を引き起こす可能性があります。
変更ストリームの制限 (該当する場合):
-
変更ストリームには MongoDB 4.0 以降が必要です。
-
変更ストリームを使用する場合、双方向同期はサポートされません。
-
非エラスティックな Amazon DocumentDB クラスターの場合は、変更ストリームを使用します。[移行方法] を [ChangeStream] に、[アーキテクチャ] を [シャーディングクラスター] に設定します。
新しいデータベース:
-
DTS は、同期タスクの開始後に作成されたデータベースからの増分データを同期しません。
サポートされる操作
キャプチャされる操作は、移行方法によって異なります。
oplog の使用 (推奨):
-
CREATE COLLECTION、CREATE INDEX -
DROP DATABASE、DROP COLLECTION、DROP INDEX -
RENAME COLLECTION -
ドキュメントレベルのINSERT、UPDATE、DELETE
変更ストリームの使用:
-
DROP DATABASE、DROP COLLECTION -
RENAME COLLECTION -
ドキュメントレベルのINSERT、UPDATE、DELETE
課金
増分データ同期は課金対象です。料金の詳細については、「課金の概要」をご参照ください。
| 課金方法 | 説明 |
|---|---|
| サブスクリプション | 1~9 か月、または 1、2、3、5 年分を前払いします。長期利用の場合は、よりコスト効率が高くなります |
| 従量課金 | 時間単位で課金されます。不要になったインスタンスをリリースすると、課金が停止します |
DTS は、接続リトライ期間中もインスタンスに課金します。
前提条件
開始する前に、以下が準備できていることを確認してください。
-
実行中の ApsaraDB for MongoDB シャーディングクラスターインスタンス。「シャーディングクラスターインスタンスの作成」をご参照ください
-
[ハンドラータイプ] が [イベントハンドラー] に設定された Function Compute のサービスと関数。「関数のクイック作成」をご参照ください
-
ソースの ApsaraDB for MongoDB インスタンス上のデータベースアカウントに、ソース、
admin、およびlocalデータベースに対する読み取り権限があること。「MongoDBデータベースユーザーの権限管理」をご参照ください説明増分同期方法として ChangeStream を使用する場合、ソースデータベースアカウントにはインスタンス全体の Change Streams 読み取り権限 (
readAnyDatabaseなど) が必要です。 ソースがカスタムアカウントを使用する ApsaraDB for MongoDB インスタンスの場合、アカウントにadminデータベースに対する読み取り権限も付与する必要があります。 詳細については、「インスタンス作成時に指定する root アカウントの権限」をご参照ください。
同期タスクの作成
ステップ1:データ同期ページの表示
いずれかのコンソールを使用して [データ同期] ページに移動します。
DTSコンソール
DMSコンソール
正確なナビゲーションパスは、DMS コンソールのレイアウトによって異なります。「シンプルモード」および「DMSコンソールのレイアウトとスタイルのカスタマイズ」をご参照ください。
ステップ2:ソースデータベースと宛先データベースの設定
[タスクの作成] をクリックし、次のパラメーターを設定します。
| セクション | パラメーター | 説明 |
|---|---|---|
| N/A | タスク名 | DTS が自動的に名前を生成します。タスクを識別しやすくするために、説明的な名前を指定します。一意である必要はありません |
| ソースデータベース | 既存の接続を選択 | 登録済みのデータベースインスタンスを選択して以下のパラメーターを自動入力するか、手動で設定します |
| データベースタイプ | [MongoDB] を選択します | |
| アクセス方法 | [Alibaba Cloudインスタンス] を選択します | |
| インスタンスリージョン | ソース MongoDB インスタンスのリージョン。 | |
| Alibaba Cloudアカウント間でのデータレプリケーション | 同一アカウントでの同期の場合は [いいえ] を選択します | |
| アーキテクチャ | [シャーディングクラスター] を選択します | |
| [移行方法] | DTS が増分データをキャプチャする方法。オプション: [Oplog] (推奨) または [ChangeStream] (詳細は後述) | |
| インスタンス ID | ソース MongoDB インスタンスの ID。 | |
| 認証データベース | アカウント認証情報を保存するデータベース。デフォルト: admin |
|
| データベースアカウント | 必要な読み取り権限を持つソースデータベースアカウント。 | |
| データベースパスワード | データベースアカウントのパスワード。 | |
| シャードアカウント | ソースインスタンスのシャードにアクセスするためのアカウント。 | |
| シャードパスワード | シャードアカウントのパスワード。 | |
| 暗号化 | 接続暗号化モード。オプション: [非暗号化]、[SSL暗号化]、または [Mongo Atlas SSL]。利用可能なオプションは、[アクセス方法] と [アーキテクチャ] の選択によって異なります。[SSL暗号化] は、[アーキテクチャ] が [シャーディングクラスター] で [移行方法] が [Oplog] の場合は利用できません | |
| 宛先データベース | 既存の接続を選択 | 登録済みのFunction Computeインスタンスを選択して以下のパラメーターを自動入力するか、手動で設定します |
| データベースタイプ | [Function Compute] を選択します | |
| アクセス方法 | [Alibaba Cloudインスタンス] を選択します | |
| インスタンスリージョン | ソースリージョンと一致します。変更はできません | |
| サービス | 宛先関数を含む Function Compute サービス。 | |
| 関数 | 同期されたデータを受信する関数。 | |
| サービスバージョンとエイリアス | サービスのバージョンまたはエイリアス。オプション: [デフォルトバージョン] (LATEST に固定)、[指定バージョン] ([サービスバージョン] が必要)、または [指定エイリアス] ([サービスエイリアス] が必要)。「用語」をご参照ください |
移行方法の選択:
| 方法 | 使用する状況 |
|---|---|
| Oplog (推奨) | ソースで oplog が有効になっている場合 (セルフマネージド MongoDB と ApsaraDB for MongoDB の両方でデフォルト)。ログのプルが高速なため、同期のレイテンシーが低くなります |
| [ChangeStream] | oplog が無効になっている場合、または非エラスティックな Amazon DocumentDB クラスターを使用している場合 (変更ストリームが必要)。MongoDB 4.0 以降が必要です。双方向同期はサポートされません |
[アーキテクチャ] が [シャーディングクラスター] で、[移行方法] が [ChangeStream] の場合、[シャードアカウント] と [シャードパスワード] のパラメーターは不要です。
ステップ3:接続テスト
[接続テストと次へ] をクリックします。
DTS サーバーの CIDR ブロックがソースと宛先の両方のセキュリティグループまたは許可リストに追加されていることを確認してください。「DTSサーバーのCIDRブロックの追加」をご参照ください。ソースまたは宛先がセルフマネージドのアクセス方法を使用している場合は、まず [接続テスト] ダイアログで [DTSサーバーのCIDRブロック] をクリックします。
ステップ4:同期オブジェクトの選択
[オブジェクトの設定] ステップで、以下を設定します。
| パラメーター | 説明 |
|---|---|
| [同期タイプ] | [増分データ同期] に固定されています。変更できません |
| [データ形式] | [Canal Json] に固定されています。フィールドの説明については、「Kafkaクラスターのトピックのデータ形式」の Canal Json セクションをご参照ください |
| [ソースオブジェクト] | 同期するデータベースまたはコレクションを選択し、 |
| [選択済みオブジェクト] | 選択したオブジェクトを確認します。 |
ステップ5:高度な設定
[次へ:高度な設定] をクリックし、以下を設定します。
| パラメーター | 説明 |
|---|---|
| [タスクスケジューリング用の専用クラスター] | デフォルトでは、DTS は共有クラスターにタスクをスケジュールします。より高い安定性を得るには、専用クラスターを購入してください。「DTS専用クラスターとは」をご参照ください |
| [接続失敗時のリトライ時間] | ソースまたは宛先に到達できない場合に DTS がリトライする時間。範囲: 10~1440 分。デフォルト: 720。少なくとも 30 分に設定してください。複数のタスクが同じソースまたは宛先を共有している場合、最短のリトライ時間が優先されます |
| [その他の問題に対するリトライ時間] | 失敗した DDL または DML 操作を DTS がリトライする時間。範囲: 1~1440 分。デフォルト: 10。少なくとも 10 分に設定してください。[接続失敗時のリトライ時間] より短くする必要があります |
| [更新後にドキュメント全体を取得] | ChangeStream のみ。[はい]:更新後に完全なドキュメントを送信します。[いいえ]:変更されたフィールドのみを送信します |
| [増分データ同期のスロットリングを有効化] | 宛先への負荷を軽減するために、同期スループットを制限します。[増分データ同期のRPS] と [増分同期のデータ同期速度 (MB/s)] を設定します |
| [環境タグ] | DTSインスタンスの環境を識別するためのタグ。オプションです |
| [ETLの設定] | ETL機能を有効にして、転送中のデータを変換します。「ETLとは」および「データ移行またはデータ同期タスクでのETLの設定」をご参照ください |
| [モニタリングとアラート] | タスクが失敗した場合、または同期のレイテンシーがしきい値を超えた場合にアラートを送信します。「DTSタスク作成時のモニタリングとアラートの設定」をご参照ください |
ステップ6:事前チェックの実行
[次へ:タスク設定を保存して事前チェック] をクリックします。
このタスク設定のAPIパラメーターをプレビューするには、[次へ:タスク設定を保存して事前チェック] にカーソルを合わせ、[OpenAPIパラメーターのプレビュー] をクリックします。
DTS は、同期タスクを開始する前に事前チェックを実行します。タスクは、事前チェックに合格した後にのみ開始されます。
-
項目が失敗した場合:失敗した項目の横にある [詳細の表示] をクリックし、報告された問題を修正してから、[再事前チェック] をクリックします。
-
アラートがトリガーされた場合:
-
アラートを無視できない場合: [詳細の表示] をクリックして問題を修正し、事前チェックを再実行します。
-
アラートを無視できる場合: [アラート詳細の確認] をクリックし、ダイアログで [無視] をクリックし、[OK] で確認してから、[再事前チェック] をクリックします。アラートを無視すると、データの不整合が発生する可能性があります。
-
ステップ7:インスタンスの購入
-
[成功率] が [100%] に達するまで待ってから、[次へ:インスタンスの購入] をクリックします。
-
[購入] ページで、以下を設定します。
| セクション | パラメーター | 説明 |
|---|---|---|
| 新しいインスタンスクラス | 課金方法 | [サブスクリプション] または [従量課金] |
| リソースグループ設定 | インスタンスのリソースグループ。デフォルト: [デフォルトリソースグループ]。「リソース管理とは」をご参照ください | |
| インスタンスクラス | 同期速度を決定します。「データ同期インスタンスのインスタンスクラス」をご参照ください | |
| サブスクリプション期間 | [サブスクリプション] 課金でのみ利用可能です。オプション: 1~9 か月、または 1、2、3、5 年 |
-
[Data Transmission Service (Pay-as-you-go) Service Terms] を読み、選択します。
-
[購入して開始] をクリックし、確認ダイアログで [OK] をクリックします。
タスクリストでタスクの進捗状況を追跡します。
次のステップ
-
Function Compute イベントハンドラーを記述して Canal JSON ペイロードを処理するには、
isDdlフィールドを使用して DDL と DML のロジックを分岐させ、各data要素のdoc文字列をデシリアライズしてドキュメントフィールドにアクセスします。 -
関数に対する大きなフィールドを持つドキュメントの負荷を軽減するには、データが関数に到達する前に ETL フィルターを設定します。「データ移行またはデータ同期タスクでのETLの設定」をご参照ください。
-
同期の健全性をモニタリングするには、[高度な設定] で、またはタスクの作成後にアラートを設定します。「DTSタスク作成時のモニタリングとアラートの設定」をご参照ください。
-
タスクが失敗した場合、DTS のテクニカルサポートが 8 時間以内に復元を試みます。復元中、タスクが再起動されたり、タスクパラメーター (データベースパラメーターではない) が調整されたりすることがあります。変更される可能性のあるパラメーターのリストについては、「インスタンスパラメーターの変更」をご参照ください。