All Products
Search
Document Center

Realtime Compute for Apache Flink:MongoDB

Last Updated:Jul 25, 2026

Konektor MongoDB mengintegrasikan ApsaraDB for MongoDB dan MongoDB yang dikelola sendiri ke Realtime Compute for Apache Flink sebagai tabel sumber, dimensi, dan sink. Konektor ini menggunakan Change Stream API untuk menangkap event insert, update, replace, dan delete secara real time.

Kemampuan

Category

Description

Jenis tabel

Sumber SQL, lookup (dimensi), dan sink

Sumber Flink CDC

Sumber DataStream

Mode berjalan

Streaming

Jenis API

DataStream API, SQL, Flink CDC

Semantik penulisan sink

Insert, update, dan delete (dengan primary key yang dideklarasikan)

Metrik pemantauan

Tabel sumber:

numBytesIn, numBytesInPerSecond, numRecordsIn, numRecordsInPerSecond, numRecordsInErrors, currentFetchEventTimeLag, currentEmitEventTimeLag, watermarkLag, sourceIdleTime

Tabel dimensi dan sink tidak mengekspos metrik pemantauan.

Untuk definisi metrik, lihat Metrics.

Cara kerja

Konektor MongoDB membaca data dalam dua fase:

  1. Snapshot penuh — membaca semua dokumen yang ada dari koleksi target secara paralel.

  2. Pembacaan inkremental — secara otomatis beralih ke mengonsumsi oplog melalui Change Stream API setelah snapshot selesai.

Proses ini menyediakan semantik tepat-sekali (exactly-once semantics), memastikan tidak ada catatan duplikat atau hilang selama pemulihan kesalahan.

Konsep utama

Startup modes

Pilih startup mode berdasarkan kapan pipeline Anda perlu mulai mengonsumsi data:

Mode

Behavior

Use when

initial (default)

Membaca snapshot saat pertama kali dijalankan, lalu beralih ke pembacaan inkremental

Anda memerlukan salinan lengkap data yang ada

latest-offset

Dimulai dari posisi oplog saat ini; tidak ada data historis

Anda hanya memerlukan perubahan mulai sekarang

timestamp

Membaca event oplog dari timestamp tertentu; melewati snapshot

Anda memerlukan perubahan mulai dari titik waktu tertentu (memerlukan MongoDB 4.0+)

Dukungan changelog lengkap

Secara default, MongoDB tidak menyimpan status dokumen sebelum perubahan (versi sebelum MongoDB 6.0). Tanpa informasi ini, konektor hanya dapat menghasilkan event UPSERT — catatan UPDATE_BEFORE tidak tersedia.

Sebagai solusi, planner Flink SQL menyisipkan operator ChangelogNormalize yang menyimpan status dokumen di state backend Flink. Meskipun berfungsi, pendekatan ini mengonsumsi penyimpanan state yang signifikan.

image.png

MongoDB 6.0+ mendukung pencatatan preimage dan postimage. Saat diaktifkan, MongoDB mencatat status dokumen lengkap sebelum dan sesudah setiap perubahan. Mengatur scan.full-changelog ke true memberi tahu konektor untuk menggunakan catatan ini guna menghasilkan aliran changelog lengkap — menghilangkan operator ChangelogNormalize dan overhead state-nya.

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

  • Instans ApsaraDB for MongoDB (replica set atau kluster sharded), atau kluster MongoDB 3.6+ yang dikelola sendiri dengan mode replica set diaktifkan. Lihat Replication.

  • Jika autentikasi diaktifkan, pengguna MongoDB dengan izin berikut: splitVector, listDatabases, listCollections, collStats, find, changeStream, serta akses baca ke config.collections dan config.chunks.

  • Alamat IP kluster Flink ditambahkan ke daftar IP yang diizinkan MongoDB.

  • Database dan koleksi target telah dibuat sebelum menjalankan pekerjaan.

Batasan

Sumber SQL

  • Pembacaan snapshot paralel memerlukan MongoDB 4.0 atau versi lebih baru. Aktifkan dengan mengatur scan.incremental.snapshot.enabled ke true.

  • Database admin, local, dan config serta semua koleksi sistem tidak dapat dipantau. Ini merupakan batasan Change Stream MongoDB. Lihat Change Streams dalam dokumentasi MongoDB.

  • Saat membuat tabel sumber SQL, deklarasikan kolom _id STRING dan tetapkan sebagai primary key.

Sink SQL

  • VVR 8.0.4 dan versi sebelumnya: hanya insert.

  • VVR 8.0.5 dan versi lebih baru dengan primary key yang dideklarasikan: insert, update, dan delete.

  • VVR 8.0.5 dan versi lebih baru tanpa primary key: hanya insert.

  • Pengiriman tepat-sekali (exactly-once delivery) tidak didukung. Opsi sink.delivery-guarantee menerima nilai none atau at-least-once.

Lookup SQL (dimensi)

  • Didukung di VVR 8.0.5 dan versi lebih baru.

  • VVR 8.0.9 dan versi lebih baru: join lookup mendukung pembacaan field bawaan _id bertipe ObjectId.

SQL

Sintaksis

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'
)
Deklarasikan kolom _id STRING dan tentukan sebagai primary key saat membuat tabel sumber CDC.

Opsi konektor

Umum

Option

Type

Required

Default

Description

connector

String

Yes

Identifier konektor. Tabel sumber: mongodb-cdc (VVR 8.0.4 dan versi sebelumnya) atau mongodb / mongodb-cdc (VVR 8.0.5 dan versi lebih baru). Tabel dimensi atau sink: mongodb.

uri

String

No

URI koneksi MongoDB. Tentukan salah satu dari uri atau hosts. Jika Anda menentukan uri, jangan sertakan scheme, hosts, username, password, dan connection.options. Jika keduanya diatur, uri memiliki prioritas lebih tinggi.

hosts

String

No

Hostname server MongoDB. Pisahkan beberapa host dengan koma (,).

scheme

String

No

mongodb

Protokol koneksi. Nilai valid: mongodb (default), mongodb+srv (DNS SRV).

username

String

No

Username MongoDB. Diperlukan saat autentikasi diaktifkan.

password

String

No

Password MongoDB. Diperlukan saat autentikasi diaktifkan. Gunakan variabel alih-alih hardcoding kredensial.

database

String

No

Nama database MongoDB. Mendukung ekspresi reguler untuk tabel sumber. Jika tidak diatur, semua database dipantau. Tidak dapat memantau database admin, local, atau config.

collection

String

No

Nama koleksi MongoDB. Mendukung ekspresi reguler untuk tabel sumber.

Catatan:

  • Jika tidak diatur, semua koleksi dipantau.

  • Tidak dapat memantau koleksi sistem.

  • Jika nama koleksi berisi karakter khusus ekspresi reguler, gunakan namespace berkualifikasi penuh (database.collection).

connection.options

String

No

Opsi koneksi tambahan sebagai pasangan key=value yang dipisahkan oleh \& (misalnya, connectTimeoutMS=12000\&socketTimeoutMS=13000).

Secara default, konektor tidak mengatur timeout koneksi socket, yang dapat menyebabkan gangguan panjang selama fluktuasi jaringan. Atur socketTimeoutMS ke nilai yang wajar untuk menghindari hal ini.

Sumber

Option

Type

Required

Default

Description

scan.startup.mode

String

No

initial

Startup mode. Nilai valid: initial, latest-offset, timestamp. Lihat Startup modes dan Startup Properties.

scan.startup.timestamp-millis

Long

Conditional

Timestamp awal dalam milidetik sejak epoch UNIX. Diperlukan saat scan.startup.mode adalah timestamp.

initial.snapshotting.queue.size

Integer

No

10240

Ukuran maksimum antrian selama fase snapshot awal. Hanya berlaku saat scan.startup.mode adalah initial.

batch.size

Integer

No

1024

Ukuran batch kursor.

poll.max.batch.size

Integer

No

1024

Jumlah maksimum dokumen perubahan yang ditarik per batch selama pembacaan stream. Nilai yang lebih besar mengalokasikan buffer internal yang lebih besar.

poll.await.time.ms

Integer

No

1000

Interval antara permintaan tarik data, dalam milidetik.

heartbeat.interval.ms

Integer

No

0

Interval heartbeat dalam milidetik. Konektor mengirim heartbeat untuk melacak posisi oplog terbaru. Mengatur nilai ini ke 0 menonaktifkan heartbeat. Atur nilai ini untuk koleksi yang jarang diperbarui.

scan.incremental.snapshot.enabled

Boolean

No

false

Mengaktifkan pembacaan snapshot paralel. Fitur eksperimental. Memerlukan MongoDB 4.0 atau versi lebih baru.

scan.incremental.snapshot.chunk.size.mb

Integer

No

64

Ukuran chunk untuk pembacaan snapshot paralel, dalam MB. Fitur eksperimental. Hanya berlaku saat pembacaan snapshot paralel diaktifkan.

scan.full-changelog

Boolean

No

false

Menghasilkan aliran changelog lengkap menggunakan catatan preimage dan postimage MongoDB. Fitur eksperimental. Memerlukan MongoDB 6.0 atau versi lebih baru dengan fitur preimage dan postimage diaktifkan.

scan.flatten-nested-columns.enabled

Boolean

No

false

Mengurai field yang dipisahkan oleh . sebagai field dokumen BSON bersarang. Misalnya, {"nested":{"col":true}} dipetakan ke field bernama nested.col. Didukung di VVR 8.0.5 dan versi lebih baru.

scan.primitive-as-string

Boolean

No

false

Mengurai semua tipe primitif BSON sebagai STRING. Didukung di VVR 8.0.5 dan versi lebih baru.

scan.ignore-delete.enabled

Boolean

No

false

Mengabaikan semua event DELETE (-D), termasuk yang dihasilkan selama pengarsipan data MongoDB. Didukung di VVR 11.1 dan versi lebih baru.

scan.incremental.snapshot.backfill.skip

Boolean

No

false

Nilai valid:

  • true: melewati backfill selama pembacaan snapshot inkremental.

  • false (default): tidak melewati backfill.

Backfill hanya berlaku selama kueri snapshot satu chunk dan tidak mencakup seluruh fase pembacaan penuh. Saat backfill dilewati, kueri snapshot setiap chunk membaca data terbaru pada saat itu; pembaruan yang terjadi pada chunk setelah dibaca tidak digabung selama fase pembacaan penuh dan dibaca dari OpLog setelah memasuki fase inkremental. Misalnya, pembaruan pada chunk5 yang terjadi saat chunk5 sedang di-snapshot tercermin langsung dalam snapshot chunk5; jika chunk5 diperbarui setelah pembaca maju ke chunk80, pembaruan tersebut diterapkan kemudian dari OpLog selama fase inkremental.

Penting

Saat diaktifkan, perubahan yang terjadi selama atau setelah pemindaian chunk tetap dikirimkan dari OpLog pada fase inkremental dan mungkin diduplikasi. Hanya semantik at-least-once yang dijamin. Aktifkan ini hanya jika sink downstream mendukung penulisan idempoten berdasarkan primary key.

Catatan

Hanya didukung di VVR 11.1 dan versi lebih baru.

initial.snapshotting.pipeline

String

No

Operasi pipeline agregasi MongoDB yang diterapkan selama pembacaan snapshot untuk memfilter data. Tentukan sebagai array JSON, misalnya: [{"$match": {"closed": "false"}}]. Hanya berlaku saat scan.startup.mode adalah initial dan konektor berjalan dalam mode Debezium. Didukung di VVR 11.1 dan versi lebih baru.

initial.snapshotting.max.threads

Integer

No

Jumlah thread untuk replikasi snapshot. Hanya berlaku saat scan.startup.mode adalah initial. Didukung di VVR 11.1 dan versi lebih baru.

initial.snapshotting.queue.size

Integer

No

16000

Ukuran antrian untuk snapshot awal. Hanya berlaku saat scan.startup.mode adalah initial. Didukung di VVR 11.1 dan versi lebih baru.

scan.change-stream.reading.parallelism

Integer

No

1

Jumlah pembaca konkuren untuk Change Stream. Hanya berlaku saat scan.incremental.snapshot.enabled adalah true. Atur juga heartbeat.interval.ms saat menggunakan opsi ini. Didukung di VVR 11.2 dan versi lebih baru.

scan.change-stream.reading.queue-size

Integer

No

16384

Ukuran antrian pesan untuk langganan change stream konkuren. Hanya berlaku saat scan.change-stream.reading.parallelism diaktifkan. Didukung di VVR 11.2 dan versi lebih baru.

Lookup (dimensi)

Option

Type

Required

Default

Description

lookup.cache

String

No

NONE

Kebijakan cache. Nilai valid: NONE (tanpa caching), PARTIAL (cache hasil lookup dari database eksternal).

lookup.max-retries

Integer

No

3

Retries maksimum saat lookup gagal.

lookup.retry.interval

Duration

No

1s

Interval antar retries saat lookup gagal.

lookup.partial-cache.expire-after-access

Duration

No

Waktu maksimum entri cache bertahan setelah akses terakhir. Unit yang didukung: ms, s, min, h, d. Memerlukan lookup.cache = PARTIAL.

lookup.partial-cache.expire-after-write

Duration

No

Waktu maksimum entri cache bertahan setelah ditulis. Memerlukan lookup.cache = PARTIAL.

lookup.partial-cache.max-rows

Long

No

Jumlah maksimum baris dalam cache. Entri terlama dihapus saat batas tercapai. Memerlukan lookup.cache = PARTIAL.

lookup.partial-cache.cache-missing-key

Boolean

No

true

Menyimpan entri null saat kunci lookup tidak memiliki catatan yang cocok. Memerlukan lookup.cache = PARTIAL.

Sink

Option

Type

Required

Default

Description

sink.buffer-flush.max-rows

Integer

No

1000

Jumlah maksimum catatan yang ditulis per batch.

sink.buffer-flush.interval

Duration

No

1s

Interval flush.

sink.delivery-guarantee

String

No

at-least-once

Semantik pengiriman tulis. Nilai valid: none, at-least-once. Exactly-once tidak didukung.

sink.max-retries

Integer

No

3

Retries maksimum saat penulisan gagal.

sink.retry.interval

Duration

No

1s

Interval antar retries saat penulisan gagal.

sink.parallelism

Integer

No

Paralelisme sink kustom.

sink.delete-strategy

String

No

CHANGELOG_STANDARD

Strategi untuk menangani event -D dan -U. Nilai valid: CHANGELOG_STANDARD (menerapkan update dan delete secara normal), IGNORE_DELETE (mengabaikan event -D; menimpa baris penuh pada -U), PARTIAL_UPDATE (mengabaikan event -U untuk mendukung update kolom parsial; menghapus baris pada -D), IGNORE_ALL (mengabaikan kedua event -U dan -D).

Pemetaan tipe data

Sumber

BSON type

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\>\>\>

Lookup (dimensi) dan sink

BSON type

Flink SQL type

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

Kolom metadata

Sumber SQL mendukung kolom metadata berikut:

Metadata column

Type

Description

database_name

STRING NOT NULL

Database yang berisi dokumen.

collection_name

STRING NOT NULL

Koleksi yang berisi dokumen.

op_ts

TIMESTAMP_LTZ(3) NOT NULL

Waktu saat dokumen berubah. Mengembalikan 0 untuk dokumen dari snapshot awal.

row_kind

STRING NOT NULL

Jenis event perubahan: +I (INSERT), -D (DELETE), -U (UPDATE_BEFORE), +U (UPDATE_AFTER). Didukung di VVR 11.1 dan versi lebih baru.

Contoh

Sumber

Contoh berikut membaca dari tabel sumber MongoDB dengan pembacaan snapshot paralel dan changelog lengkap diaktifkan, lalu menulis field yang dipilih ke sink print.

-- Tabel sumber CDC: membaca data produk dari MongoDB
-- _id harus dideklarasikan dan ditetapkan sebagai primary key
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', -- Aktifkan pembacaan snapshot paralel (memerlukan MongoDB 4.0+)
  'scan.full-changelog' = 'true'                -- Aktifkan changelog lengkap (memerlukan MongoDB 6.0+ dengan preimage/postimage)
);

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;

Lookup (dimensi)

Contoh berikut melakukan join antara stream generator data dengan tabel dimensi MongoDB menggunakan temporal join.

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',                           -- Cache hasil lookup untuk performa lebih baik
  'lookup.partial-cache.expire-after-access' = '10min', -- Hapus entri cache setelah 10 menit tidak aktif
  'lookup.partial-cache.expire-after-write' = '10min',  -- Hapus entri cache 10 menit setelah ditulis
  'lookup.partial-cache.max-rows' = '100'               -- Maksimal 100 baris dalam cache
);

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;

Sink

Contoh berikut menulis data dari generator data ke tabel sink MongoDB. Primary key dideklarasikan untuk mendukung operasi insert, update, dan delete.

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 -- Deklarasikan primary key untuk mengaktifkan update dan delete
) 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 (pratinjau publik)

Flink CDC memungkinkan Anda menyinkronkan data MongoDB ke penyimpanan downstream menggunakan pipeline berbasis skrip YAML, tanpa menulis DDL SQL. Fitur ini memerlukan VVR 11.1 dan versi lebih baru.

Sintaksis

source:
   type: mongodb
   name: MongoDB Source
   hosts: localhost:33076
   username: ${mongo.username}
   password: ${mongo.password}
   database: foo_db
   collection: foo_col_.*

sink:
  type: ...

Opsi konfigurasi

Option

Required

Type

Default

Description

type

Yes

STRING

Konektor. Atur ke mongodb.

scheme

No

STRING

mongodb

Protokol koneksi. Nilai valid: mongodb, mongodb+srv.

hosts

Yes

STRING

Hostname server MongoDB. Pisahkan beberapa host dengan koma.

username

No

STRING

Username MongoDB.

password

No

STRING

Password MongoDB.

database

Yes

STRING

Nama database MongoDB yang akan ditangkap. Ekspresi reguler didukung.

collection

Yes

STRING

Nama koleksi MongoDB yang akan ditangkap. Ekspresi reguler didukung. Gunakan namespace berkualifikasi penuh database.collection.

connection.options

No

STRING

Opsi koneksi tambahan sebagai pasangan k=v yang dipisahkan oleh \&, misalnya: replicaSet=test\&connectTimeoutMS=300000.

schema.inference.strategy

No

STRING

continuous

Strategi inferensi skema. continuous: menginferensi tipe secara kontinu dan mengeluarkan event perubahan skema saat skema melebar. static: menginferensi skema sekali saat startup.

scan.max.pre.fetch.records

No

INT

50

Jumlah maksimum catatan yang disampel per koleksi selama inferensi skema awal.

scan.startup.mode

No

STRING

initial

Startup mode. Nilai valid: initial, latest-offset, timestamp, snapshot.

scan.startup.timestamp-millis

No

LONG

Timestamp awal dalam milidetik. Diperlukan saat scan.startup.mode adalah timestamp.

chunk-meta.group.size

No

INT

1000

Ukuran maksimum chunk metadata.

scan.incremental.close-idle-reader.enabled

No

BOOLEAN

false

Menutup pembaca sumber yang idle setelah beralih ke pembacaan inkremental.

scan.incremental.snapshot.backfill.skip

No

BOOLEAN

false

Nilai valid:

  • true: melewati backfill selama pembacaan snapshot inkremental.

  • false (default): tidak melewati backfill.

Backfill hanya berlaku selama kueri snapshot satu chunk dan tidak mencakup seluruh fase pembacaan penuh. Saat backfill dilewati, kueri snapshot setiap chunk membaca data terbaru pada saat itu; pembaruan yang terjadi pada chunk setelah dibaca tidak digabung selama fase pembacaan penuh dan dibaca dari OpLog setelah memasuki fase inkremental. Misalnya, pembaruan pada chunk5 yang terjadi saat chunk5 sedang di-snapshot tercermin langsung dalam snapshot chunk5; jika chunk5 diperbarui setelah pembaca maju ke chunk80, pembaruan tersebut diterapkan kemudian dari OpLog selama fase inkremental.

Penting

Saat diaktifkan, perubahan yang terjadi selama atau setelah pemindaian chunk tetap dikirimkan dari OpLog pada fase inkremental dan mungkin diduplikasi. Hanya semantik at-least-once yang dijamin. Aktifkan ini hanya jika sink downstream mendukung penulisan idempoten berdasarkan primary key.

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

No

BOOLEAN

false

Membaca chunk tak terbatas terlebih dahulu. Mengurangi risiko kehabisan memori untuk koleksi yang sering diperbarui.

batch.size

No

INT

1024

Ukuran batch kursor.

poll.max.batch.size

No

INT

1024

Jumlah maksimum entri per permintaan tarik Change Stream.

poll.await.time.ms

No

INT

1000

Waktu tunggu minimum antara permintaan tarik Change Stream, dalam milidetik.

heartbeat.interval.ms

No

INT

0

Interval heartbeat dalam milidetik. Atur ini untuk koleksi yang jarang diperbarui. Mengatur ke 0 menonaktifkan heartbeat.

scan.incremental.snapshot.chunk.size.mb

No

INT

64

Ukuran chunk selama snapshotting, dalam MB.

scan.incremental.snapshot.chunk.samples

No

INT

20

Jumlah sampel yang digunakan untuk memperkirakan ukuran koleksi selama snapshotting.

scan.full-changelog

No

BOOLEAN

false

Menghasilkan event changelog lengkap menggunakan catatan preimage dan postimage. Memerlukan MongoDB 6.0 atau versi lebih baru dengan preimage dan postimage diaktifkan.

scan.cursor.no-timeout

No

BOOLEAN

false

Menonaktifkan timeout kursor. Secara default, MongoDB menutup kursor idle setelah 10 menit.

scan.ignore-delete.enabled

No

BOOLEAN

false

Mengabaikan event delete dari MongoDB.

scan.flatten.nested-documents.enabled

No

BOOLEAN

false

Meratakan dokumen BSON bersarang. Misalnya, {"doc": {"foo": 1, "bar": "two"}} menjadi doc.foo INT, doc.bar STRING.

scan.all.primitives.as-string.enabled

No

BOOLEAN

false

Menginferensi semua tipe primitif sebagai STRING. Mengurangi event perubahan skema saat tipe upstream tidak konsisten.

metadata.list

No

STRING

Daftar field metadata yang dipisahkan koma untuk dikirimkan ke downstream. Nilai yang didukung: ts_ms (timestamp event OpLog), op_ts (alias untuk ts_ms; gunakan saat menulis metadata ke Kafka JSON).

Pemetaan tipe data

MongoDB BSON

Flink CDC

Notes

STRING

VARCHAR

INT32

INT

INT64

BIGINT

DECIMAL128

DECIMAL

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

TIMESTAMP

TIMESTAMP

DATETIME

LOCALZONEDTIMESTAMP

BINARY

VARBINARY

DOCUMENT

MAP

Tipe kunci dan nilai diinferensi.

ARRAY

ARRAY

Tipe elemen diinferensi.

OBJECTID

VARCHAR

Direpresentasikan sebagai string heksadesimal.

SYMBOL, REGULAREXPRESSION, JAVASCRIPT, JAVASCRIPTWITHSCOPE

VARCHAR

Direpresentasikan sebagai string.

Kolom metadata

Flink CDC mendukung kolom metadata berikut untuk konektor MongoDB:

Metadata column

Type

Description

ts_ms

BIGINT NOT NULL

Waktu saat dokumen berubah (timestamp OpLog). Mengembalikan 0 untuk dokumen dari snapshot awal.

Gunakan kolom metadata generik modul Transform untuk mengakses database_name, collection_name, dan row_kind.

DataStream API

Penting

Untuk menggunakan DataStream API, siapkan konektor DataStream untuk pekerjaan Anda. Lihat DataStream connector usage.

Tambahkan dependensi Maven

Repositori Maven Central menyediakan konektor MongoDB VVR.

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>flink-connector-mongodb</artifactId>
    <version>${vvr.version}</version>
</dependency>

Buat MongoDBSource

Gunakan MongoDBSource.builder() untuk membuat sumber:

  • Untuk mengaktifkan pembacaan snapshot inkremental, gunakan builder dari com.ververica.cdc.connectors.mongodb.source.

  • Jika tidak, gunakan builder dari com.ververica.cdc.connectors.mongodb.

MongoDBSource.builder()
  .hosts("mongo.example.com:27017")
  .username("mongouser")
  .password("mongopasswd")
  .databaseList("testdb")          // Mendukung ekspresi reguler; gunakan .* untuk mencocokkan semua database
  .collectionList("testcoll")      // Mendukung ekspresi reguler; gunakan .* untuk mencocokkan semua koleksi
  .startupOptions(StartupOptions.initial()) // StartupOptions.latest-offset(), StartupOptions.timestamp()
  .deserializer(new JsonDebeziumDeserializationSchema())
  .build();

Parameter MongoDBSource

Parameter

Description

hosts

Hostname server MongoDB.

username

Username MongoDB. Abaikan jika autentikasi tidak diaktifkan.

password

Password MongoDB. Abaikan jika autentikasi tidak diaktifkan.

databaseList

Nama database yang akan dipantau. Mendukung ekspresi reguler. Gunakan .* untuk mencocokkan semua database.

collectionList

Nama koleksi yang akan dipantau. Mendukung ekspresi reguler. Gunakan .* untuk mencocokkan semua koleksi.

startupOptions

Mode startup. Nilai yang valid: StartupOptions.initial(), StartupOptions.latest-offset(), dan StartupOptions.timestamp().

deserializer

Deserializer untuk mengonversi objek SourceRecord. Nilai valid: MongoDBConnectorDeserializationSchema (mode upsert, menghasilkan RowData Flink), MongoDBConnectorFullChangelogDeserializationSchema (mode changelog lengkap, menghasilkan RowData Flink), JsonDebeziumDeserializationSchema (menghasilkan string JSON).

Referensi