MongoDB コネクタは、ApsaraDB for MongoDB およびセルフマネージド MongoDB を、ソーステーブル、ディメンションテーブル、結果テーブルとして Realtime Compute for Apache Flink に統合します。このコネクタは Change Stream API を使用して、挿入、更新、置換、削除の各イベントをリアルタイムでキャプチャします。
機能
|
カテゴリ |
説明 |
|
テーブルタイプ |
SQL ソース、ルックアップ (ディメンション)、およびシンク Flink CDC ソース DataStream ソース |
|
実行モード |
ストリーミング |
|
API タイプ |
DataStream API、SQL、Flink CDC |
|
シンクの書き込みセマンティクス |
挿入、更新、削除 (プライマリキーが宣言されている場合) |
監視メトリクス
ソーステーブル:
numBytesIn、numBytesInPerSecond、numRecordsIn、numRecordsInPerSecond、numRecordsInErrors、currentFetchEventTimeLag、currentEmitEventTimeLag、watermarkLag、sourceIdleTime
ディメンションテーブルと結果テーブルは監視メトリクスを公開しません。
メトリクスの定義については、「メトリクス」をご参照ください。
仕組み
MongoDB コネクタは、2 つのフェーズでデータを読み取ります:
-
フルスナップショット:対象のコレクションから既存のすべてのドキュメントを並行して読み取ります。
-
増分読み取り:スナップショットが完了すると、自動的に Change Stream API を介した oplog の消費に切り替わります。
このプロセスは exactly-once セマンティクスを提供し、障害復旧中にレコードの重複や欠落がないことを保証します。
基本概念
起動モード
パイプラインがデータの消費を開始するタイミングに基づいて起動モードを選択します:
|
モード |
動作 |
使用シーン |
|
|
初回起動時にスナップショットを読み取り、その後増分読み取りに切り替えます |
既存のデータの完全なコピーが必要です |
|
|
現在の oplog の位置から開始します。既存データは読み取りません |
現時点以降の変更のみが必要な場合 |
|
|
指定されたタイムスタンプから oplog イベントを読み取ります。スナップショットはスキップします |
既知の時点からの変更が必要な場合 (MongoDB 4.0 以降が必要) |
フルチェンジログのサポート
デフォルトでは、MongoDB はドキュメントの変更前の状態を保存しません (MongoDB 6.0 より前のバージョン)。この情報がないと、コネクタは UPSERT イベントしか生成できず、UPDATE_BEFORE レコードが欠落します。
この問題を回避するため、Flink SQL プランナーは ChangelogNormalize 演算子を挿入し、ドキュメントの状態を Flink の状態バックエンドにキャッシュします。このアプローチは機能しますが、大量の状態ストレージを消費します。

MongoDB 6.0 以降は、preimage と postimage の記録をサポートしています。有効にすると、MongoDB は各変更の前後で完全なドキュメント状態を記録します。scan.full-changelog を true に設定すると、コネクタはこれらのレコードを使用してフルチェンジログストリームを生成するように指示されます。これにより、ChangelogNormalize 演算子とその状態オーバーヘッドが不要になります。
前提条件
開始する前に、以下が準備できていることを確認してください:
-
ApsaraDB for MongoDB インスタンス (レプリカセットまたはシャードクラスター)、またはレプリカセットモードが有効なセルフマネージド MongoDB 3.6 以降のクラスター。詳細については、「レプリケーション」をご参照ください。
-
認証が有効な場合、次の権限を持つ MongoDB ユーザー:
splitVector、listDatabases、listCollections、collStats、find、changeStream、およびconfig.collectionsとconfig.chunksへの読み取りアクセス。 -
MongoDB の IP 許可リストに追加された Flink クラスターの IP アドレス。
-
ジョブを実行する前に作成されたターゲットデータベースとコレクション。
制限事項
SQL ソース
-
並列スナップショット読み取りには MongoDB 4.0 以降が必要です。
scan.incremental.snapshot.enabledをtrueに設定して有効にします。 -
admin、local、configデータベースおよびすべてのシステムコレクションは監視できません。これは MongoDB Change Stream の制限です。MongoDB ドキュメントの「Change Streams」をご参照ください。 -
SQL ソーステーブルを作成する際、
_id STRING列を宣言し、それをプライマリキーとして設定してください。
SQL シンク
-
VVR 8.0.4 以前:挿入のみ。
-
VVR 8.0.5 以降でプライマリキーが宣言されている場合:挿入、更新、削除。
-
VVR 8.0.5 以降でプライマリキーがない場合:挿入のみ。
-
exactly-once 配信はサポートされていません。
sink.delivery-guaranteeオプションはnoneまたはat-least-onceを受け入れます。
SQL ルックアップ (ディメンション)
-
VVR 8.0.5 以降でサポートされています。
-
VVR 8.0.9 以降:ルックアップ結合は、ObjectId 型の組み込み
_idフィールドの読み取りをサポートします。
SQL
構文
CREATE TABLE tableName(
_id STRING,
[columnName dataType,]*
PRIMARY KEY(_id) NOT ENFORCED
) WITH (
'connector' = 'mongodb',
'hosts' = 'localhost:27017',
'username' = 'mongouser',
'password' = '${secret_values.password}',
'database' = 'testdb',
'collection' = 'testcoll'
)
CDC ソーステーブルを作成する際は、_id STRING 列を宣言し、それをプライマリキーとして指定してください。
コネクタオプション
一般
|
オプション |
型 |
必須 |
デフォルト |
説明 |
|
|
String |
はい |
— |
コネクタ識別子。ソーステーブル: |
|
|
String |
いいえ |
— |
MongoDB 接続 URI。 |
|
|
String |
いいえ |
— |
MongoDB サーバーのホスト名。複数のホストはカンマ ( |
|
|
String |
いいえ |
|
接続プロトコル。有効な値: |
|
|
String |
いいえ |
— |
MongoDB のユーザー名。認証が有効な場合に必須です。 |
|
|
String |
いいえ |
— |
MongoDB のパスワード。認証が有効な場合に必須です。認証情報をハードコーディングする代わりに変数を使用してください。 |
|
|
String |
いいえ |
— |
MongoDB のデータベース名。ソーステーブルでは正規表現をサポートします。設定しない場合、すべてのデータベースが監視されます。 |
|
|
String |
いいえ |
— |
MongoDB のコレクション名。ソーステーブルでは正規表現をサポートします。 注意:
|
|
|
String |
いいえ |
— |
追加の接続オプションを デフォルトでは、コネクタはソケット接続タイムアウトを設定しないため、ネットワークジッター中に長時間の中断が発生する可能性があります。これを避けるために、 |
ソース
|
オプション |
型 |
必須 |
デフォルト |
説明 |
|
|
String |
いいえ |
|
起動モード。有効な値: |
|
|
Long |
条件付き |
— |
UNIX エポックからの開始タイムスタンプ (ミリ秒単位)。 |
|
|
Integer |
いいえ |
|
初期スナップショットフェーズ中の最大キューサイズ。 |
|
|
Integer |
いいえ |
|
カーソルのバッチサイズ。 |
|
|
Integer |
いいえ |
|
ストリーム読み取り中にバッチごとにプルされる変更ドキュメントの最大数。値が大きいほど、より大きな内部バッファーが割り当てられます。 |
|
|
Integer |
いいえ |
|
データプルリクエスト間の間隔 (ミリ秒単位)。 |
|
|
Integer |
いいえ |
|
ハートビート間隔 (ミリ秒単位)。コネクタはハートビートを送信して最新の oplog 位置を追跡します。これを |
|
|
Boolean |
いいえ |
|
並列スナップショット読み取りを有効にします。実験的な機能です。MongoDB 4.0 以降が必要です。 |
|
|
Integer |
いいえ |
|
並列スナップショット読み取りのチャンクサイズ (MB 単位)。実験的な機能です。並列スナップショット読み取りが有効な場合にのみ有効です。 |
|
|
Boolean |
いいえ |
|
MongoDB の preimage と postimage レコードを使用してフルチェンジログストリームを生成します。実験的な機能です。preimage と postimage 機能が有効な MongoDB 6.0 以降が必要です。 |
|
|
Boolean |
いいえ |
|
|
|
|
Boolean |
いいえ |
|
すべてのプリミティブ BSON 型を STRING として解析します。VVR 8.0.5 以降でサポートされています。 |
|
|
Boolean |
いいえ |
|
MongoDB のデータアーカイブ中に生成されるものを含む、すべての DELETE (-D) イベントを無視します。VVR 11.1 以降でサポートされています。 |
|
|
Boolean |
いいえ |
|
有効な値:
バックフィルは単一チャンクのスナップショットクエリ中にのみ適用され、フル読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその瞬間の最新データを読み取ります。チャンクが読み取られた後に発生した更新は、フル読み取りフェーズ中にはマージされず、増分フェーズに入った後に OpLog から読み取られます。たとえば、chunk5 のスナップショット中に chunk5 への更新が発生した場合、その更新は chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に OpLog から適用されます。 重要
有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで OpLog から配信され、重複する可能性があります。at-least-once セマンティクスのみが保証されます。ダウンストリームの結果テーブルがプライマリキーによるべき等書き込みをサポートしている場合にのみ有効にしてください。 説明
VVR 11.1 以降でのみサポートされています。 |
|
|
String |
いいえ |
— |
スナップショット読み取り中にデータをフィルターするために適用される MongoDB 集約パイプライン操作。JSON 配列として指定します。例: |
|
|
Integer |
いいえ |
— |
スナップショットレプリケーション用のスレッド数。 |
|
|
Integer |
いいえ |
|
初期スナップショット用のキューサイズ。 |
|
|
Integer |
いいえ |
|
Change Stream の同時リーダー数。 |
|
|
Integer |
いいえ |
|
同時 Change Stream サブスクリプション用のメッセージキューサイズ。 |
ルックアップ (ディメンション)
|
オプション |
型 |
必須 |
デフォルト |
説明 |
|
|
String |
いいえ |
|
キャッシュポリシー。有効な値: |
|
|
Integer |
いいえ |
|
ルックアップ失敗時の最大リトライ回数。 |
|
|
Duration |
いいえ |
|
ルックアップ失敗時のリトライ間隔。 |
|
|
Duration |
いいえ |
— |
キャッシュされたエントリが最終アクセス後に生存する最大時間。サポートされる単位: |
|
|
Duration |
いいえ |
— |
キャッシュされたエントリが書き込まれた後に生存する最大時間。 |
|
|
Long |
いいえ |
— |
キャッシュ内の最大行数。制限に達すると、最も古いエントリが削除されます。 |
|
|
Boolean |
いいえ |
|
ルックアップキーに一致するレコードがない場合に null エントリをキャッシュします。 |
|
String |
いいえ |
|
Flink の文字列値が MongoDB ディメンションテーブルと照合される際に使用される変換戦略。有効な値:
説明 このパラメーターは VVR 11.9 以降でのみサポートされています。 |
シンク
|
オプション |
型 |
必須 |
デフォルト |
説明 |
|
|
Integer |
いいえ |
|
バッチごとに書き込まれるレコードの最大数。 |
|
|
Duration |
いいえ |
|
フラッシュ間隔。 |
|
|
String |
いいえ |
|
書き込み配信セマンティクス。有効な値: |
|
|
Integer |
いいえ |
|
書き込み失敗時の最大リトライ回数。 |
|
|
Duration |
いいえ |
|
書き込み失敗時のリトライ間隔。 |
|
|
Integer |
いいえ |
— |
カスタムシンクの並列度。 |
|
|
String |
いいえ |
|
-D および -U イベントの処理戦略。有効な値: |
データ型マッピング
ソース
|
BSON 型 |
Flink SQL |
|
Int32 |
INT |
|
Int64 |
BIGINT |
|
Double |
DOUBLE |
|
Decimal128 |
DECIMAL(p, s) |
|
Boolean |
BOOLEAN |
|
Date Timestamp |
DATE |
|
Date Timestamp |
TIME |
|
DateTime |
TIMESTAMP(3), TIMESTAMP_LTZ(3) |
|
Timestamp |
TIMESTAMP(0), TIMESTAMP_LTZ(0) |
|
String, ObjectId, UUID, Symbol, MD5, JavaScript, Regex |
STRING |
|
Binary |
BYTES |
|
Object |
ROW |
|
Array |
ARRAY |
|
DBPointer |
ROW\<$ref STRING, $id STRING\> |
|
GeoJSON Point |
ROW\<type STRING, coordinates ARRAY\<DOUBLE\>\> |
|
GeoJSON Line |
ROW\<type STRING, coordinates ARRAY\<ARRAY\<DOUBLE\>\>\> |
ルックアップ (ディメンション) と結果テーブル
|
BSON 型 |
Flink SQL 型 |
|
Int32 |
INT |
|
Int64 |
BIGINT |
|
Double |
DOUBLE |
|
Decimal128 |
DECIMAL |
|
Boolean |
BOOLEAN |
|
DateTime |
TIMESTAMP_LTZ(3) |
|
Timestamp |
TIMESTAMP_LTZ(0) |
|
String, ObjectId |
STRING |
|
Binary |
BYTES |
|
Object |
ROW |
|
Array |
ARRAY |
メタデータ列
SQL ソースは以下のメタデータ列をサポートします:
|
メタデータ列 |
型 |
説明 |
|
|
STRING NOT NULL |
ドキュメントを含むデータベース。 |
|
|
STRING NOT NULL |
ドキュメントを含むコレクション。 |
|
|
TIMESTAMP_LTZ(3) NOT NULL |
ドキュメントが変更された時刻。初期スナップショットからのドキュメントの場合は |
|
|
STRING NOT NULL |
変更イベントタイプ: |
例
ソース
以下の例では、並列スナップショット読み取りとフルチェンジログを有効にした MongoDB ソーステーブルから読み取り、選択したフィールドを print 結果テーブルに書き込みます。
-- CDC ソーステーブル: MongoDB からプロダクトデータを読み取る
-- _id を宣言し、プライマリキーとして設定する必要がある
CREATE TEMPORARY TABLE mongo_source (
`_id` STRING,
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price ROW<amount DECIMAL, currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>,
db_name STRING METADATA FROM 'database_name' VIRTUAL,
collection_name STRING METADATA VIRTUAL,
op_ts TIMESTAMP_LTZ(3) METADATA VIRTUAL,
PRIMARY KEY(_id) NOT ENFORCED
) WITH (
'connector' = 'mongodb',
'hosts' = 'dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,',
'username' = 'root',
'password' = '${secret_values.password}',
'database' = 'flinktest',
'collection' = 'flinkcollection',
'scan.incremental.snapshot.enabled' = 'true', -- 並列スナップショット読み取りを有効にする (MongoDB 4.0+ が必要)
'scan.full-changelog' = 'true' -- フルチェンジログを有効にする (preimage/postimage が有効な MongoDB 6.0+ が必要)
);
CREATE TEMPORARY TABLE productssink (
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price_amount DECIMAL,
suppliers_name STRING,
db_name STRING,
collection_name STRING,
op_ts TIMESTAMP_LTZ(3)
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO productssink
SELECT
name,
weight,
tags,
price.amount,
suppliers[1].name,
db_name,
collection_name,
op_ts
FROM mongo_source;
ルックアップ (ディメンション)
以下の例では、データジェネレーターストリームを時間的な結合を使用して MongoDB ディメンションテーブルと結合します。
CREATE TEMPORARY TABLE datagen_source (
id STRING,
a INT,
b BIGINT,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE mongo_dim (
`_id` STRING,
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price ROW<amount DECIMAL, currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>,
PRIMARY KEY(_id) NOT ENFORCED
) WITH (
'connector' = 'mongodb',
'hosts' = 'dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,',
'username' = 'root',
'password' = '${secret_values.password}',
'database' = 'flinktest',
'collection' = 'flinkcollection',
'lookup.cache' = 'PARTIAL', -- パフォーマンス向上のためルックアップ結果をキャッシュする
'lookup.partial-cache.expire-after-access' = '10min', -- 10分間非アクティブなキャッシュエントリを削除する
'lookup.partial-cache.expire-after-write' = '10min', -- 書き込み後10分でキャッシュエントリを削除する
'lookup.partial-cache.max-rows' = '100' -- キャッシュ内の最大行数を100に設定
);
CREATE TEMPORARY TABLE print_sink (
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price_amount DECIMAL,
suppliers_name STRING
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO print_sink
SELECT
T.id,
T.a,
T.b,
H.name
FROM datagen_source AS T
JOIN mongo_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H ON T.id = H._id;
シンク
以下の例では、データジェネレーターから MongoDB 結果テーブルにデータを書き込みます。挿入、更新、削除操作をサポートするためにプライマリキーが宣言されています。
CREATE TEMPORARY TABLE datagen_source (
`_id` STRING,
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price ROW<amount DECIMAL, currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE mongo_sink (
`_id` STRING,
name STRING,
weight DECIMAL,
tags ARRAY<STRING>,
price ROW<amount DECIMAL, currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>,
PRIMARY KEY(_id) NOT ENFORCED -- 更新と削除を有効にするためにプライマリキーを宣言
) WITH (
'connector' = 'mongodb',
'hosts' = 'dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,dds-bp169b982fc25****.mongodb.rds.aliyuncs.com:3717,',
'username' = 'root',
'password' = '${secret_values.password}',
'database' = 'flinktest',
'collection' = 'flinkcollection'
);
INSERT INTO mongo_sink SELECT * FROM datagen_source;
Flink CDC (パブリックプレビュー)
Flink CDC を使用すると、SQL DDL を記述することなく、YAML スクリプトベースのパイプラインを使用して MongoDB データをダウンストリームストアに同期できます。この機能には VVR 11.1 以降が必要です。
構文
source:
type: mongodb
name: MongoDB Source
hosts: localhost:33076
username: ${mongo.username}
password: ${mongo.password}
database: foo_db
collection: foo_col_.*
sink:
type: ...
構成オプション
|
オプション |
必須 |
型 |
デフォルト |
説明 |
|
|
はい |
STRING |
— |
コネクタ。 |
|
|
いいえ |
STRING |
|
接続プロトコル。有効な値: |
|
|
はい |
STRING |
— |
MongoDB サーバーのホスト名。複数のホストはカンマで区切ります。 |
|
|
いいえ |
STRING |
— |
MongoDB のユーザー名。 |
|
|
いいえ |
STRING |
— |
MongoDB のパスワード。 |
|
|
はい |
STRING |
— |
キャプチャする MongoDB データベース名。正規表現がサポートされています。 |
|
|
はい |
STRING |
— |
キャプチャする MongoDB コレクション名。正規表現がサポートされています。完全修飾 |
|
|
いいえ |
STRING |
— |
追加の接続オプションを |
|
|
いいえ |
STRING |
|
スキーマ推論戦略。 |
|
|
いいえ |
INT |
|
初期スキーマ推論中にコレクションごとにサンプリングする最大レコード数。 |
|
|
いいえ |
STRING |
|
起動モード。有効な値: |
|
|
いいえ |
LONG |
— |
開始タイムスタンプ (ミリ秒単位)。 |
|
|
いいえ |
INT |
|
最大メタデータチャンクサイズ。 |
|
|
いいえ |
BOOLEAN |
|
増分読み取りに切り替えた後、アイドル状態のソースリーダーを閉じます。 |
|
|
いいえ |
BOOLEAN |
|
有効な値:
バックフィルは単一チャンクのスナップショットクエリ中にのみ適用され、フル読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその瞬間の最新データを読み取ります。チャンクが読み取られた後に発生した更新は、フル読み取りフェーズ中にはマージされず、増分フェーズに入った後に OpLog から読み取られます。たとえば、chunk5 のスナップショット中に chunk5 への更新が発生した場合、その更新は chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に OpLog から適用されます。 重要
有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで OpLog から配信され、重複する可能性があります。at-least-once セマンティクスのみが保証されます。ダウンストリームの結果テーブルがプライマリキーによるべき等書き込みをサポートしている場合にのみ有効にしてください。 |
|
|
いいえ |
BOOLEAN |
|
有界でないチャンクを最初に読み取ります。更新頻度の高いコレクションのメモリ不足リスクを軽減します。 |
|
|
いいえ |
INT |
|
カーソルのバッチサイズ。 |
|
|
いいえ |
INT |
|
Change Stream プルリクエストごとの最大エントリ数。 |
|
|
いいえ |
INT |
|
Change Stream プルリクエスト間の最小待機時間 (ミリ秒単位)。 |
|
|
いいえ |
INT |
|
ハートビート間隔 (ミリ秒単位)。更新頻度の低いコレクションに対してこれを設定してください。 |
|
|
いいえ |
INT |
|
スナップショット中のチャンクサイズ (MB 単位)。 |
|
|
いいえ |
INT |
|
スナップショット中にコレクションサイズを推定するために使用されるサンプル数。 |
|
|
いいえ |
BOOLEAN |
|
preimage と postimage レコードを使用してフルチェンジログイベントを生成します。preimage と postimage が有効な MongoDB 6.0 以降が必要です。 |
|
|
いいえ |
BOOLEAN |
|
カーソルのタイムアウトを無効にします。デフォルトでは、MongoDB は 10 分後にアイドル状態のカーソルを閉じます。 |
|
|
いいえ |
BOOLEAN |
|
MongoDB からの削除イベントを無視します。 |
|
|
いいえ |
BOOLEAN |
|
ネストされた BSON ドキュメントをフラット化します。たとえば、 |
|
|
いいえ |
BOOLEAN |
|
すべてのプリミティブ型を STRING として推論します。アップストリームの型が一致しない場合にスキーマ変更イベントを減らします。 |
|
|
いいえ |
STRING |
— |
ダウンストリームに渡すメタデータフィールドのカンマ区切りリスト。サポートされる値: |
データ型マッピング
|
MongoDB BSON |
Flink CDC |
注意 |
|
STRING |
VARCHAR |
— |
|
INT32 |
INT |
— |
|
INT64 |
BIGINT |
— |
|
DECIMAL128 |
DECIMAL |
— |
|
DOUBLE |
DOUBLE |
— |
|
BOOLEAN |
BOOLEAN |
— |
|
TIMESTAMP |
TIMESTAMP |
— |
|
DATETIME |
LOCALZONEDTIMESTAMP |
— |
|
BINARY |
VARBINARY |
— |
|
DOCUMENT |
MAP |
キーと値の型が推論されます。 |
|
ARRAY |
ARRAY |
要素の型が推論されます。 |
|
OBJECTID |
VARCHAR |
16 進数文字列として表現されます。 |
|
SYMBOL, REGULAREXPRESSION, JAVASCRIPT, JAVASCRIPTWITHSCOPE |
VARCHAR |
文字列として表現されます。 |
メタデータ列
Flink CDC は MongoDB コネクタに対して以下のメタデータ列をサポートします:
|
メタデータ列 |
型 |
説明 |
|
|
BIGINT NOT NULL |
ドキュメントが変更された時刻 (OpLog タイムスタンプ)。初期スナップショットからのドキュメントの場合は |
Transform モジュールの汎用メタデータ列を使用して、database_name、collection_name、および row_kind にアクセスします。
DataStream API
DataStream API を使用するには、ジョブ用に DataStream コネクタを設定する必要があります。詳細については、「DataStream コネクタの使用方法」をご参照ください。
Maven 依存関係の追加
Maven Central Repository は VVR MongoDB コネクタをホストしています。
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>flink-connector-mongodb</artifactId>
<version>${vvr.version}</version>
</dependency>
MongoDBSource のビルド
MongoDBSource.builder() を使用してソースを構築します:
-
増分スナップショット読み取りを有効にするには、
com.ververica.cdc.connectors.mongodb.sourceのビルダーを使用します。 -
それ以外の場合は、
com.ververica.cdc.connectors.mongodbのビルダーを使用します。
MongoDBSource.builder()
.hosts("mongo.example.com:27017")
.username("mongouser")
.password("mongopasswd")
.databaseList("testdb") // 正規表現をサポート。すべてのデータベースに一致させるには .* を使用
.collectionList("testcoll") // 正規表現をサポート。すべてのコレクションに一致させるには .* を使用
.startupOptions(StartupOptions.initial()) // StartupOptions.latest-offset(), StartupOptions.timestamp()
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
MongoDBSource パラメーター
|
パラメーター |
説明 |
|
|
MongoDB サーバーのホスト名。 |
|
|
MongoDB のユーザー名。認証が有効でない場合は省略します。 |
|
|
MongoDB のパスワード。認証が有効でない場合は省略します。 |
|
|
監視するデータベース名。正規表現をサポートします。すべてのデータベースに一致させるには |
|
|
監視するコレクション名。正規表現をサポートします。すべてのコレクションに一致させるには |
|
|
起動モード。有効な値: |
|
|
|
リファレンス
-
Flink CDC (パブリックプレビュー) — MongoDB のデータとスキーマの変更をダウンストリームテーブルに同期します (VVR 11.1 以降)。
-
メトリクス — ソーステーブルのパフォーマンスを監視します。