すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:MongoDB

最終更新日:Sep 19, 2026

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 つのフェーズでデータを読み取ります:

  1. フルスナップショット:対象のコレクションから既存のすべてのドキュメントを並行して読み取ります。

  2. 増分読み取り:スナップショットが完了すると、自動的に Change Stream API を介した oplog の消費に切り替わります。

このプロセスは exactly-once セマンティクスを提供し、障害復旧中にレコードの重複や欠落がないことを保証します。

基本概念

起動モード

パイプラインがデータの消費を開始するタイミングに基づいて起動モードを選択します:

モード

動作

使用シーン

initial (デフォルト)

初回起動時にスナップショットを読み取り、その後増分読み取りに切り替えます

既存のデータの完全なコピーが必要です

latest-offset

現在の oplog の位置から開始します。既存データは読み取りません

現時点以降の変更のみが必要な場合

timestamp

指定されたタイムスタンプから oplog イベントを読み取ります。スナップショットはスキップします

既知の時点からの変更が必要な場合 (MongoDB 4.0 以降が必要)

フルチェンジログのサポート

デフォルトでは、MongoDB はドキュメントの変更前の状態を保存しません (MongoDB 6.0 より前のバージョン)。この情報がないと、コネクタは UPSERT イベントしか生成できず、UPDATE_BEFORE レコードが欠落します。

この問題を回避するため、Flink SQL プランナーは ChangelogNormalize 演算子を挿入し、ドキュメントの状態を Flink の状態バックエンドにキャッシュします。このアプローチは機能しますが、大量の状態ストレージを消費します。

image.png

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 列を宣言し、それをプライマリキーとして指定してください。

コネクタオプション

一般

オプション

型

必須

デフォルト

説明

connector

String

はい

—

コネクタ識別子。ソーステーブル:mongodb-cdc (VVR 8.0.4 以前) または mongodb / mongodb-cdc (VVR 8.0.5 以降)。ディメンションテーブルまたは結果テーブル:mongodb。

uri

String

いいえ

—

MongoDB 接続 URI。uri または hosts のいずれかを指定します。uri を指定した場合、scheme、hosts、username、password、および connection.options は省略します。両方が設定されている場合、uri が優先されます。

hosts

String

いいえ

—

MongoDB サーバーのホスト名。複数のホストはカンマ (,) で区切ります。

scheme

String

いいえ

mongodb

接続プロトコル。有効な値:mongodb (デフォルト)、mongodb+srv (DNS SRV)。

username

String

いいえ

—

MongoDB のユーザー名。認証が有効な場合に必須です。

password

String

いいえ

—

MongoDB のパスワード。認証が有効な場合に必須です。認証情報をハードコーディングする代わりに変数を使用してください。

database

String

いいえ

—

MongoDB のデータベース名。ソーステーブルでは正規表現をサポートします。設定しない場合、すべてのデータベースが監視されます。admin、local、または config データベースは監視できません。

collection

String

いいえ

—

MongoDB のコレクション名。ソーステーブルでは正規表現をサポートします。

注意:

  • 設定しない場合、すべてのコレクションが監視されます。

  • システムコレクションは監視できません。

  • コレクション名に正規表現の特殊文字が含まれる場合は、完全修飾名前空間 (database.collection) を使用してください。

connection.options

String

いいえ

—

追加の接続オプションを & で区切られた key=value のペアとして指定します (例:connectTimeoutMS=12000&socketTimeoutMS=13000)。

デフォルトでは、コネクタはソケット接続タイムアウトを設定しないため、ネットワークジッター中に長時間の中断が発生する可能性があります。これを避けるために、socketTimeoutMS に適切な値を設定してください。

ソース

オプション

型

必須

デフォルト

説明

scan.startup.mode

String

いいえ

initial

起動モード。有効な値:initial、latest-offset、timestamp。詳細については、「起動モード」および「Startup Properties」をご参照ください。

scan.startup.timestamp-millis

Long

条件付き

—

UNIX エポックからの開始タイムスタンプ (ミリ秒単位)。scan.startup.mode が timestamp の場合に必須です。

initial.snapshotting.queue.size

Integer

いいえ

10240

初期スナップショットフェーズ中の最大キューサイズ。scan.startup.mode が initial の場合にのみ有効です。

batch.size

Integer

いいえ

1024

カーソルのバッチサイズ。

poll.max.batch.size

Integer

いいえ

1024

ストリーム読み取り中にバッチごとにプルされる変更ドキュメントの最大数。値が大きいほど、より大きな内部バッファーが割り当てられます。

poll.await.time.ms

Integer

いいえ

1000

データプルリクエスト間の間隔 (ミリ秒単位)。

heartbeat.interval.ms

Integer

いいえ

0

ハートビート間隔 (ミリ秒単位)。コネクタはハートビートを送信して最新の oplog 位置を追跡します。これを 0 に設定するとハートビートが無効になります。更新頻度の低いコレクションに対してこれを設定してください。

scan.incremental.snapshot.enabled

Boolean

いいえ

false

並列スナップショット読み取りを有効にします。実験的な機能です。MongoDB 4.0 以降が必要です。

scan.incremental.snapshot.chunk.size.mb

Integer

いいえ

64

並列スナップショット読み取りのチャンクサイズ (MB 単位)。実験的な機能です。並列スナップショット読み取りが有効な場合にのみ有効です。

scan.full-changelog

Boolean

いいえ

false

MongoDB の preimage と postimage レコードを使用してフルチェンジログストリームを生成します。実験的な機能です。preimage と postimage 機能が有効な MongoDB 6.0 以降が必要です。

scan.flatten-nested-columns.enabled

Boolean

いいえ

false

. で区切られたフィールドをネストされた BSON ドキュメントフィールドとして解析します。たとえば、{"nested":{"col":true}} は nested.col という名前のフィールドにマッピングされます。VVR 8.0.5 以降でサポートされています。

scan.primitive-as-string

Boolean

いいえ

false

すべてのプリミティブ BSON 型を STRING として解析します。VVR 8.0.5 以降でサポートされています。

scan.ignore-delete.enabled

Boolean

いいえ

false

MongoDB のデータアーカイブ中に生成されるものを含む、すべての DELETE (-D) イベントを無視します。VVR 11.1 以降でサポートされています。

scan.incremental.snapshot.backfill.skip

Boolean

いいえ

false

有効な値:

  • true:増分スナップショット読み取り中にバックフィルをスキップします。

  • false (デフォルト):バックフィルをスキップしません。

バックフィルは単一チャンクのスナップショットクエリ中にのみ適用され、フル読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその瞬間の最新データを読み取ります。チャンクが読み取られた後に発生した更新は、フル読み取りフェーズ中にはマージされず、増分フェーズに入った後に OpLog から読み取られます。たとえば、chunk5 のスナップショット中に chunk5 への更新が発生した場合、その更新は chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に OpLog から適用されます。

重要

有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで OpLog から配信され、重複する可能性があります。at-least-once セマンティクスのみが保証されます。ダウンストリームの結果テーブルがプライマリキーによるべき等書き込みをサポートしている場合にのみ有効にしてください。

説明

VVR 11.1 以降でのみサポートされています。

initial.snapshotting.pipeline

String

いいえ

—

スナップショット読み取り中にデータをフィルターするために適用される MongoDB 集約パイプライン操作。JSON 配列として指定します。例:[{"$match": {"closed": "false"}}]。scan.startup.mode が initial で、コネクタが Debezium モードで実行される場合にのみ有効です。VVR 11.1 以降でサポートされています。

initial.snapshotting.max.threads

Integer

いいえ

—

スナップショットレプリケーション用のスレッド数。scan.startup.mode が initial の場合にのみ有効です。VVR 11.1 以降でサポートされています。

initial.snapshotting.queue.size

Integer

いいえ

16000

初期スナップショット用のキューサイズ。scan.startup.mode が initial の場合にのみ有効です。VVR 11.1 以降でサポートされています。

scan.change-stream.reading.parallelism

Integer

いいえ

1

Change Stream の同時リーダー数。scan.incremental.snapshot.enabled が true の場合にのみ有効です。このオプションを使用する場合は heartbeat.interval.ms も設定してください。VVR 11.2 以降でサポートされています。

scan.change-stream.reading.queue-size

Integer

いいえ

16384

同時 Change Stream サブスクリプション用のメッセージキューサイズ。scan.change-stream.reading.parallelism が有効な場合にのみ有効です。VVR 11.2 以降でサポートされています。

ルックアップ (ディメンション)

オプション

型

必須

デフォルト

説明

lookup.cache

String

いいえ

NONE

キャッシュポリシー。有効な値:NONE (キャッシュなし)、PARTIAL (外部データベースからのルックアップ結果をキャッシュ)。

lookup.max-retries

Integer

いいえ

3

ルックアップ失敗時の最大リトライ回数。

lookup.retry.interval

Duration

いいえ

1s

ルックアップ失敗時のリトライ間隔。

lookup.partial-cache.expire-after-access

Duration

いいえ

—

キャッシュされたエントリが最終アクセス後に生存する最大時間。サポートされる単位:ms、s、min、h、d。lookup.cache = PARTIAL が必要です。

lookup.partial-cache.expire-after-write

Duration

いいえ

—

キャッシュされたエントリが書き込まれた後に生存する最大時間。lookup.cache = PARTIAL が必要です。

lookup.partial-cache.max-rows

Long

いいえ

—

キャッシュ内の最大行数。制限に達すると、最も古いエントリが削除されます。lookup.cache = PARTIAL が必要です。

lookup.partial-cache.cache-missing-key

Boolean

いいえ

true

ルックアップキーに一致するレコードがない場合に null エントリをキャッシュします。lookup.cache = PARTIAL が必要です。

lookup.type-conversion.mode

String

いいえ

FORCE_STRING

Flink の文字列値が MongoDB ディメンションテーブルと照合される際に使用される変換戦略。有効な値:

  • FORCE_STRING:照合前に常に MongoDB フィールドを文字列に変換します。ObjectId のように Flink に相当する型がないフィールド型も照合できますが、フィールド上のインデックス (もしあれば) は使用できません。

  • OBJECT_ID_WRAPPER:Flink からの文字列値が有効な ObjectId の 16 進数文字列である場合、それを ObjectId として照合します。結合フィールドのインデックスを使用してクエリを高速化できます。

  • NONE:特別な処理を行いません。MongoDB の型が Flink のクエリ型と異なる場合、行は照合されません。

説明

このパラメーターは VVR 11.9 以降でのみサポートされています。

シンク

オプション

型

必須

デフォルト

説明

sink.buffer-flush.max-rows

Integer

いいえ

1000

バッチごとに書き込まれるレコードの最大数。

sink.buffer-flush.interval

Duration

いいえ

1s

フラッシュ間隔。

sink.delivery-guarantee

String

いいえ

at-least-once

書き込み配信セマンティクス。有効な値:none、at-least-once。exactly-once はサポートされていません。

sink.max-retries

Integer

いいえ

3

書き込み失敗時の最大リトライ回数。

sink.retry.interval

Duration

いいえ

1s

書き込み失敗時のリトライ間隔。

sink.parallelism

Integer

いいえ

—

カスタムシンクの並列度。

sink.delete-strategy

String

いいえ

CHANGELOG_STANDARD

-D および -U イベントの処理戦略。有効な値:CHANGELOG_STANDARD (更新と削除を通常通り適用)、IGNORE_DELETE (-D イベントを無視し、-U で行全体を上書き)、PARTIAL_UPDATE (部分的な列更新をサポートするために -U イベントを無視し、-D で行を削除)、IGNORE_ALL (-U と -D の両方のイベントを無視)。

データ型マッピング

ソース

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 ソースは以下のメタデータ列をサポートします:

メタデータ列

型

説明

database_name

STRING NOT NULL

ドキュメントを含むデータベース。

collection_name

STRING NOT NULL

ドキュメントを含むコレクション。

op_ts

TIMESTAMP_LTZ(3) NOT NULL

ドキュメントが変更された時刻。初期スナップショットからのドキュメントの場合は 0 を返します。

row_kind

STRING NOT NULL

変更イベントタイプ:+I (INSERT)、-D (DELETE)、-U (UPDATE_BEFORE)、+U (UPDATE_AFTER)。VVR 11.1 以降でサポートされています。

例

ソース

以下の例では、並列スナップショット読み取りとフルチェンジログを有効にした 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: ...

構成オプション

オプション

必須

型

デフォルト

説明

type

はい

STRING

—

コネクタ。mongodb に設定します。

scheme

いいえ

STRING

mongodb

接続プロトコル。有効な値:mongodb、mongodb+srv。

hosts

はい

STRING

—

MongoDB サーバーのホスト名。複数のホストはカンマで区切ります。

username

いいえ

STRING

—

MongoDB のユーザー名。

password

いいえ

STRING

—

MongoDB のパスワード。

database

はい

STRING

—

キャプチャする MongoDB データベース名。正規表現がサポートされています。

collection

はい

STRING

—

キャプチャする MongoDB コレクション名。正規表現がサポートされています。完全修飾 database.collection 名前空間を使用してください。

connection.options

いいえ

STRING

—

追加の接続オプションを & で区切られた k=v のペアとして指定します。例:replicaSet=test&connectTimeoutMS=300000。

schema.inference.strategy

いいえ

STRING

continuous

スキーマ推論戦略。continuous:継続的に型を推論し、スキーマが拡張されたときにスキーマ変更イベントを発行します。static:起動時に一度だけスキーマを推論します。

scan.max.pre.fetch.records

いいえ

INT

50

初期スキーマ推論中にコレクションごとにサンプリングする最大レコード数。

scan.startup.mode

いいえ

STRING

initial

起動モード。有効な値:initial、latest-offset、timestamp、snapshot。

scan.startup.timestamp-millis

いいえ

LONG

—

開始タイムスタンプ (ミリ秒単位)。scan.startup.mode が timestamp の場合に必須です。

chunk-meta.group.size

いいえ

INT

1000

最大メタデータチャンクサイズ。

scan.incremental.close-idle-reader.enabled

いいえ

BOOLEAN

false

増分読み取りに切り替えた後、アイドル状態のソースリーダーを閉じます。

scan.incremental.snapshot.backfill.skip

いいえ

BOOLEAN

false

有効な値:

  • true:増分スナップショット読み取り中にバックフィルをスキップします。

  • false (デフォルト):バックフィルをスキップしません。

バックフィルは単一チャンクのスナップショットクエリ中にのみ適用され、フル読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその瞬間の最新データを読み取ります。チャンクが読み取られた後に発生した更新は、フル読み取りフェーズ中にはマージされず、増分フェーズに入った後に OpLog から読み取られます。たとえば、chunk5 のスナップショット中に chunk5 への更新が発生した場合、その更新は chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に OpLog から適用されます。

重要

有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで OpLog から配信され、重複する可能性があります。at-least-once セマンティクスのみが保証されます。ダウンストリームの結果テーブルがプライマリキーによるべき等書き込みをサポートしている場合にのみ有効にしてください。

scan.incremental.snapshot.unbounded-chunk-first.enabled

いいえ

BOOLEAN

false

有界でないチャンクを最初に読み取ります。更新頻度の高いコレクションのメモリ不足リスクを軽減します。

batch.size

いいえ

INT

1024

カーソルのバッチサイズ。

poll.max.batch.size

いいえ

INT

1024

Change Stream プルリクエストごとの最大エントリ数。

poll.await.time.ms

いいえ

INT

1000

Change Stream プルリクエスト間の最小待機時間 (ミリ秒単位)。

heartbeat.interval.ms

いいえ

INT

0

ハートビート間隔 (ミリ秒単位)。更新頻度の低いコレクションに対してこれを設定してください。0 に設定するとハートビートが無効になります。

scan.incremental.snapshot.chunk.size.mb

いいえ

INT

64

スナップショット中のチャンクサイズ (MB 単位)。

scan.incremental.snapshot.chunk.samples

いいえ

INT

20

スナップショット中にコレクションサイズを推定するために使用されるサンプル数。

scan.full-changelog

いいえ

BOOLEAN

false

preimage と postimage レコードを使用してフルチェンジログイベントを生成します。preimage と postimage が有効な MongoDB 6.0 以降が必要です。

scan.cursor.no-timeout

いいえ

BOOLEAN

false

カーソルのタイムアウトを無効にします。デフォルトでは、MongoDB は 10 分後にアイドル状態のカーソルを閉じます。

scan.ignore-delete.enabled

いいえ

BOOLEAN

false

MongoDB からの削除イベントを無視します。

scan.flatten.nested-documents.enabled

いいえ

BOOLEAN

false

ネストされた BSON ドキュメントをフラット化します。たとえば、{"doc": {"foo": 1, "bar": "two"}} は doc.foo INT, doc.bar STRING になります。

scan.all.primitives.as-string.enabled

いいえ

BOOLEAN

false

すべてのプリミティブ型を STRING として推論します。アップストリームの型が一致しない場合にスキーマ変更イベントを減らします。

metadata.list

いいえ

STRING

—

ダウンストリームに渡すメタデータフィールドのカンマ区切りリスト。サポートされる値:ts_ms (OpLog イベントタイムスタンプ)、op_ts (ts_ms のエイリアス。メタデータを Kafka JSON に書き込む際に使用)。

データ型マッピング

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 コネクタに対して以下のメタデータ列をサポートします:

メタデータ列

型

説明

ts_ms

BIGINT NOT NULL

ドキュメントが変更された時刻 (OpLog タイムスタンプ)。初期スナップショットからのドキュメントの場合は 0 を返します。

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 パラメーター

パラメーター

説明

hosts

MongoDB サーバーのホスト名。

username

MongoDB のユーザー名。認証が有効でない場合は省略します。

password

MongoDB のパスワード。認証が有効でない場合は省略します。

databaseList

監視するデータベース名。正規表現をサポートします。すべてのデータベースに一致させるには .* を使用します。

collectionList

監視するコレクション名。正規表現をサポートします。すべてのコレクションに一致させるには .* を使用します。

startupOptions

起動モード。有効な値:StartupOptions.initial()、StartupOptions.latest-offset()、StartupOptions.timestamp()。

deserializer

SourceRecord オブジェクトを変換するためのデシリアライザー。有効な値:MongoDBConnectorDeserializationSchema (アップサートモード、Flink RowData を生成)、MongoDBConnectorFullChangelogDeserializationSchema (フルチェンジログモード、Flink RowData を生成)、JsonDebeziumDeserializationSchema (JSON 文字列を生成)。

リファレンス