Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:MongoDB

Dernière mise à jour :Aug 12, 2026

Le connecteur MongoDB intègre ApsaraDB for MongoDB et les instances MongoDB auto-gérées à Realtime Compute for Apache Flink en tant que tables source, de dimension et sink. Il s'appuie sur l'API Change Stream pour capturer en temps réel les événements d'insertion, de mise à jour, de remplacement et de suppression.

Fonctionnalités

Catégorie

Description

Types de tables

Source SQL, lookup (dimension) et sink

Source Flink CDC

Source DataStream

Mode d'exécution

Streaming

Types d'API

DataStream API, SQL, Flink CDC

Sémantique d'écriture du sink

Insertion, mise à jour et suppression (avec clé primaire déclarée)

Métriques de surveillance

Tables source :

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

Les tables de dimension et sink n'exposent aucune métrique de surveillance.

Pour la définition des métriques, consultez Métriques.

Fonctionnement

Le connecteur MongoDB lit les données en deux phases :

  1. Snapshot complet : lecture parallèle de tous les documents existants dans les collections ciblées.

  2. Lecture incrémentielle : basculement automatique vers la consommation de l'oplog via l'API Change Stream une fois le snapshot terminé.

Ce processus garantit une sémantique exactly-once, assurant l'absence de doublons ou de pertes d'enregistrements lors de la reprise après incident.

Concepts clés

Modes de démarrage

Choisissez un mode de démarrage selon le moment où votre pipeline doit commencer à consommer les données :

Mode

Comportement

Cas d'usage

initial (par défaut)

Lit un snapshot au premier démarrage, puis passe en lecture incrémentielle

Une copie complète des données existantes est nécessaire

latest-offset

Démarre à partir de la position actuelle de l'oplog ; aucune donnée historique

Seules les modifications futures sont requises

timestamp

Lit les événements de l'oplog à partir d'un horodatage spécifié ; ignore le snapshot

Les modifications doivent être capturées depuis un point précis dans le temps (nécessite MongoDB 4.0+)

Prise en charge du journal de modifications complet

Par défaut, MongoDB ne conserve pas l'état antérieur des documents avant modification (versions antérieures à MongoDB 6.0). En l'absence de cette information, le connecteur ne peut produire que des événements UPSERT : les enregistrements UPDATE_BEFORE sont manquants.

Pour contourner cette limitation, le planificateur Flink SQL insère un opérateur ChangelogNormalize qui met en cache l'état des documents dans le backend d'état de Flink. Bien que fonctionnelle, cette approche consomme un espace de stockage d'état important.

image.png

MongoDB 6.0+ prend en charge l'enregistrement des images avant et après modification (preimage et postimage). Lorsque cette option est activée, MongoDB enregistre l'état complet du document avant et après chaque changement. Définir scan.full-changelog sur true indique au connecteur d'utiliser ces enregistrements pour produire des flux de journal de modifications complets, éliminant ainsi l'opérateur ChangelogNormalize et la surcharge d'état associée.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Une instance ApsaraDB for MongoDB (replica set ou cluster shardé), ou un cluster MongoDB auto-géré version 3.6+ avec le mode replica set activé. Consultez Réplication.

  • Si l'authentification est activée, un utilisateur MongoDB disposant des permissions suivantes : splitVector, listDatabases, listCollections, collStats, find, changeStream, ainsi que d'un accès en lecture à config.collections et config.chunks.

  • Les adresses IP du cluster Flink ajoutées à la liste d'autorisation IP de MongoDB.

  • La base de données et la collection cibles créées avant l'exécution du job.

Limitations

Source SQL

  • La lecture parallèle de snapshot nécessite MongoDB 4.0 ou ultérieur. Activez-la en définissant scan.incremental.snapshot.enabled sur true.

  • Les bases de données admin, local et config, ainsi que toutes les collections système, ne peuvent pas être surveillées. Cette restriction provient de l'API Change Stream de MongoDB. Consultez la documentation Change Streams de MongoDB.

  • Lors de la création d'une table source SQL, déclarez la colonne _id STRING et définissez-la comme clé primaire.

Sink SQL

  • VVR 8.0.4 et versions antérieures : insertion uniquement.

  • VVR 8.0.5 et versions ultérieures avec clé primaire déclarée : insertion, mise à jour et suppression.

  • VVR 8.0.5 et versions ultérieures sans clé primaire : insertion uniquement.

  • La livraison exactly-once n'est pas prise en charge. L'option sink.delivery-guarantee accepte les valeurs none ou at-least-once.

Lookup SQL (dimension)

  • Pris en charge dans VVR 8.0.5 et versions ultérieures.

  • VVR 8.0.9 et versions ultérieures : les jointures lookup permettent de lire le champ intégré _id de type ObjectId.

SQL

Syntaxe

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'
)
Déclarez la colonne _id STRING et spécifiez-la comme clé primaire lors de la création d'une table source CDC.

Options du connecteur

Général

Option

Type

Obligatoire

Valeur par défaut

Description

connector

String

Oui

Identifiant du connecteur. Tables source : mongodb-cdc (VVR 8.0.4 et antérieur) ou mongodb / mongodb-cdc (VVR 8.0.5 et ultérieur). Tables de dimension ou sink : mongodb.

uri

String

Non

URI de connexion MongoDB. Spécifiez soit uri, soit hosts. Si vous spécifiez uri, omettez scheme, hosts, username, password et connection.options. Si les deux sont définis, uri est prioritaire.

hosts

String

Non

Nom d'hôte du serveur MongoDB. Séparez plusieurs hôtes par des virgules (,).

scheme

String

Non

mongodb

Protocole de connexion. Valeurs valides : mongodb (par défaut), mongodb+srv (DNS SRV).

username

String

Non

Nom d'utilisateur MongoDB. Obligatoire lorsque l'authentification est activée.

password

String

Non

Mot de passe MongoDB. Obligatoire lorsque l'authentification est activée. Utilisez des variables plutôt que de coder les identifiants en dur.

database

String

Non

Nom de la base de données MongoDB. Prend en charge les expressions régulières pour les tables source. Si non défini, toutes les bases de données sont surveillées. Les bases de données admin, local et config ne peuvent pas être surveillées.

collection

String

Non

Nom de la collection MongoDB. Prend en charge les expressions régulières pour les tables source.

Notes :

  • Si non défini, toutes les collections sont surveillées.

  • Les collections système ne peuvent pas être surveillées.

  • Si le nom de la collection contient des caractères spéciaux d'expression régulière, utilisez le namespace pleinement qualifié (database.collection).

connection.options

String

Non

Options de connexion supplémentaires sous forme de paires key=value séparées par \& (par exemple, connectTimeoutMS=12000\&socketTimeoutMS=13000).

Par défaut, le connecteur ne définit pas de délai d'expiration de connexion socket, ce qui peut entraîner de longues interruptions en cas d'instabilité réseau. Définissez socketTimeoutMS sur une valeur raisonnable pour éviter cela.

Source

Option

Type

Obligatoire

Valeur par défaut

Description

scan.startup.mode

String

Non

initial

Mode de démarrage. Valeurs valides : initial, latest-offset, timestamp. Consultez Modes de démarrage et Startup Properties.

scan.startup.timestamp-millis

Long

Conditionnel

Horodatage de début en millisecondes depuis l'époque UNIX. Obligatoire lorsque scan.startup.mode est défini sur timestamp.

initial.snapshotting.queue.size

Integer

Non

10240

Taille maximale de la file d'attente pendant la phase de snapshot initial. Ne prend effet que lorsque scan.startup.mode est défini sur initial.

batch.size

Integer

Non

1024

Taille du lot du curseur.

poll.max.batch.size

Integer

Non

1024

Nombre maximal de documents de modification récupérés par lot lors de la lecture du flux. Des valeurs plus élevées allouent un tampon interne plus grand.

poll.await.time.ms

Integer

Non

1000

Intervalle entre les requêtes de récupération de données, en millisecondes.

heartbeat.interval.ms

Integer

Non

0

Intervalle de battement de cœur en millisecondes. Le connecteur envoie des battements de cœur pour suivre la dernière position de l'oplog. Définir cette valeur sur 0 désactive les battements de cœur. Configurez cette option pour les collections rarement mises à jour.

scan.incremental.snapshot.enabled

Boolean

Non

false

Active la lecture parallèle de snapshot. Fonctionnalité expérimentale. Nécessite MongoDB 4.0 ou ultérieur.

scan.incremental.snapshot.chunk.size.mb

Integer

Non

64

Taille des chunks pour la lecture parallèle de snapshot, en Mo. Fonctionnalité expérimentale. Ne prend effet que lorsque la lecture parallèle de snapshot est activée.

scan.full-changelog

Boolean

Non

false

Génère un flux de journal de modifications complet à l'aide des enregistrements preimage et postimage de MongoDB. Fonctionnalité expérimentale. Nécessite MongoDB 6.0 ou ultérieur avec les fonctionnalités preimage et postimage activées.

scan.flatten-nested-columns.enabled

Boolean

Non

false

Analyse les champs séparés par . comme des champs de document BSON imbriqué. Par exemple, {"nested":{"col":true}} est mappé à un champ nommé nested.col. Pris en charge dans VVR 8.0.5 et versions ultérieures.

scan.primitive-as-string

Boolean

Non

false

Analyse tous les types BSON primitifs comme STRING. Pris en charge dans VVR 8.0.5 et versions ultérieures.

scan.ignore-delete.enabled

Boolean

Non

false

Ignore tous les événements DELETE (-D), y compris ceux générés lors de l'archivage des données MongoDB. Pris en charge dans VVR 11.1 et versions ultérieures.

scan.incremental.snapshot.backfill.skip

Boolean

Non

false

Valeurs valides :

  • true : ignore le backfill lors de la lecture incrémentielle du snapshot.

  • false (par défaut) : n'ignore pas le backfill.

Le backfill s'applique uniquement lors de la requête de snapshot d'un seul chunk et ne couvre pas l'ensemble de la phase de lecture complète. Lorsque le backfill est ignoré, la requête de snapshot de chaque chunk lit les données les plus récentes à cet instant ; les mises à jour qui surviennent sur un chunk après sa lecture ne sont pas fusionnées pendant la phase de lecture complète et sont lues depuis l'OpLog après l'entrée en phase incrémentielle. Par exemple, une mise à jour de chunk5 survenant pendant le snapshot de chunk5 est reflétée directement dans le snapshot de chunk5 ; si chunk5 est mis à jour après que le lecteur est passé à chunk80, la mise à jour est appliquée ultérieurement depuis l'OpLog lors de la phase incrémentielle.

Important

Lorsque cette option est activée, les modifications survenant pendant ou après l'analyse d'un chunk sont toujours livrées depuis l'OpLog en phase incrémentielle et peuvent être dupliquées. Seule la sémantique at-least-once est garantie. N'activez cette option que si le sink en aval prend en charge les écritures idempotentes par clé primaire.

Remarque

Pris en charge uniquement dans VVR 11.1 et versions ultérieures.

initial.snapshotting.pipeline

String

Non

Opérations de pipeline d'agrégation MongoDB appliquées lors de la lecture du snapshot pour filtrer les données. Spécifiez un tableau JSON, par exemple : [{"$match": {"closed": "false"}}]. Ne prend effet que lorsque scan.startup.mode est défini sur initial et que le connecteur s'exécute en mode Debezium. Pris en charge dans VVR 11.1 et versions ultérieures.

initial.snapshotting.max.threads

Integer

Non

Nombre de threads pour la réplication du snapshot. Ne prend effet que lorsque scan.startup.mode est défini sur initial. Pris en charge dans VVR 11.1 et versions ultérieures.

initial.snapshotting.queue.size

Integer

Non

16000

Taille de la file d'attente pour le snapshot initial. Ne prend effet que lorsque scan.startup.mode est défini sur initial. Pris en charge dans VVR 11.1 et versions ultérieures.

scan.change-stream.reading.parallelism

Integer

Non

1

Nombre de lecteurs concurrents pour le Change Stream. Ne prend effet que lorsque scan.incremental.snapshot.enabled est défini sur true. Définissez également heartbeat.interval.ms lors de l'utilisation de cette option. Pris en charge dans VVR 11.2 et versions ultérieures.

scan.change-stream.reading.queue-size

Integer

Non

16384

Taille de la file d'attente des messages pour les abonnements concurrents au change stream. Ne prend effet que lorsque scan.change-stream.reading.parallelism est activé. Pris en charge dans VVR 11.2 et versions ultérieures.

Lookup (dimension)

Option

Type

Obligatoire

Valeur par défaut

Description

lookup.cache

String

Non

NONE

Politique de cache. Valeurs valides : NONE (aucun cache), PARTIAL (mise en cache des résultats de lookup depuis la base de données externe).

lookup.max-retries

Integer

Non

3

Nombre maximal de tentatives en cas d'échec du lookup.

lookup.retry.interval

Duration

Non

1s

Intervalle entre les tentatives en cas d'échec du lookup.

lookup.partial-cache.expire-after-access

Duration

Non

Durée de vie maximale d'une entrée en cache après le dernier accès. Unités prises en charge : ms, s, min, h, d. Nécessite lookup.cache = PARTIAL.

lookup.partial-cache.expire-after-write

Duration

Non

Durée de vie maximale d'une entrée en cache après son écriture. Nécessite lookup.cache = PARTIAL.

lookup.partial-cache.max-rows

Long

Non

Nombre maximal de lignes dans le cache. Les entrées les plus anciennes sont évictées lorsque la limite est atteinte. Nécessite lookup.cache = PARTIAL.

lookup.partial-cache.cache-missing-key

Boolean

Non

true

Met en cache une entrée nulle lorsqu'une clé de lookup ne correspond à aucun enregistrement. Nécessite lookup.cache = PARTIAL.

Sink

Option

Type

Obligatoire

Valeur par défaut

Description

sink.buffer-flush.max-rows

Integer

Non

1000

Nombre maximal d'enregistrements écrits par lot.

sink.buffer-flush.interval

Duration

Non

1s

Intervalle de vidage du tampon.

sink.delivery-guarantee

String

Non

at-least-once

Sémantique de livraison des écritures. Valeurs valides : none, at-least-once. La sémantique exactly-once n'est pas prise en charge.

sink.max-retries

Integer

Non

3

Nombre maximal de tentatives en cas d'échec d'écriture.

sink.retry.interval

Duration

Non

1s

Intervalle entre les tentatives en cas d'échec d'écriture.

sink.parallelism

Integer

Non

Parallélisme personnalisé du sink.

sink.delete-strategy

String

Non

CHANGELOG_STANDARD

Stratégie de gestion des événements -D et -U. Valeurs valides : CHANGELOG_STANDARD (applique normalement les mises à jour et suppressions), IGNORE_DELETE (ignore les événements -D ; écrase les lignes complètes sur -U), PARTIAL_UPDATE (ignore les événements -U pour prendre en charge les mises à jour partielles de colonnes ; supprime les lignes sur -D), IGNORE_ALL (ignore à la fois les événements -U et -D).

Mappages des types de données

Source

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

Lookup (dimension) et sink

Type BSON

Type 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

Colonnes de métadonnées

La source SQL prend en charge les colonnes de métadonnées suivantes :

Colonne de métadonnées

Type

Description

database_name

STRING NOT NULL

Base de données contenant le document.

collection_name

STRING NOT NULL

Collection contenant le document.

op_ts

TIMESTAMP_LTZ(3) NOT NULL

Heure de modification du document. Retourne 0 pour les documents issus du snapshot initial.

row_kind

STRING NOT NULL

Type d'événement de modification : +I (INSERT), -D (DELETE), -U (UPDATE_BEFORE), +U (UPDATE_AFTER). Pris en charge dans VVR 11.1 et versions ultérieures.

Exemples

Source

L'exemple suivant lit les données d'une table source MongoDB avec lecture parallèle de snapshot et journal de modifications complet activés, puis écrit les champs sélectionnés dans un sink print.

-- CDC source table: reads product data from MongoDB
-- _id must be declared and set as the 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', -- Enable parallel snapshot reading (requires MongoDB 4.0+)
  'scan.full-changelog' = 'true'                -- Enable full changelog (requires MongoDB 6.0+ with 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 (dimension)

L'exemple suivant joint un flux de générateur de données à une table de dimension MongoDB à l'aide d'une jointure temporelle.

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 lookup results for better performance
  'lookup.partial-cache.expire-after-access' = '10min', -- Evict cached entries after 10 minutes of inactivity
  'lookup.partial-cache.expire-after-write' = '10min',  -- Evict cached entries 10 minutes after write
  'lookup.partial-cache.max-rows' = '100'               -- Maximum 100 rows in the 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

L'exemple suivant écrit les données d'un générateur de données dans une table sink MongoDB. Une clé primaire est déclarée pour permettre les opérations d'insertion, de mise à jour et de suppression.

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 -- Declare primary key to enable update and 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 (aperçu public)

Flink CDC permet de synchroniser les données MongoDB vers un stockage en aval à l'aide d'un pipeline basé sur un script YAML, sans écrire de DDL SQL. Cette fonctionnalité nécessite VVR 11.1 ou ultérieur.

Syntaxe

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

sink:
  type: ...

Options de configuration

Option

Obligatoire

Type

Valeur par défaut

Description

type

Oui

STRING

Connecteur. Définissez sur mongodb.

scheme

Non

STRING

mongodb

Protocole de connexion. Valeurs valides : mongodb, mongodb+srv.

hosts

Oui

STRING

Nom(s) d'hôte du serveur MongoDB. Séparez plusieurs hôtes par des virgules.

username

Non

STRING

Nom d'utilisateur MongoDB.

password

Non

STRING

Mot de passe MongoDB.

database

Oui

STRING

Nom de la base de données MongoDB à capturer. Les expressions régulières sont prises en charge.

collection

Oui

STRING

Nom de la collection MongoDB à capturer. Les expressions régulières sont prises en charge. Utilisez le namespace pleinement qualifié database.collection.

connection.options

Non

STRING

Options de connexion supplémentaires sous forme de paires k=v séparées par \&, par exemple : replicaSet=test\&connectTimeoutMS=300000.

schema.inference.strategy

Non

STRING

continuous

Stratégie d'inférence de schéma. continuous : infère les types en continu et émet des événements de changement de schéma lorsque le schéma s'élargit. static : infère le schéma une seule fois au démarrage.

scan.max.pre.fetch.records

Non

INT

50

Nombre maximal d'enregistrements à échantillonner par collection lors de l'inférence initiale du schéma.

scan.startup.mode

Non

STRING

initial

Mode de démarrage. Valeurs valides : initial, latest-offset, timestamp, snapshot.

scan.startup.timestamp-millis

Non

LONG

Horodatage de début en millisecondes. Obligatoire lorsque scan.startup.mode est défini sur timestamp.

chunk-meta.group.size

Non

INT

1000

Taille maximale des chunks de métadonnées.

scan.incremental.close-idle-reader.enabled

Non

BOOLEAN

false

Ferme les lecteurs source inactifs après le basculement en lecture incrémentielle.

scan.incremental.snapshot.backfill.skip

Non

BOOLEAN

false

Valeurs valides :

  • true : ignore le backfill lors de la lecture incrémentielle du snapshot.

  • false (par défaut) : n'ignore pas le backfill.

Le backfill s'applique uniquement lors de la requête de snapshot d'un seul chunk et ne couvre pas l'ensemble de la phase de lecture complète. Lorsque le backfill est ignoré, la requête de snapshot de chaque chunk lit les données les plus récentes à cet instant ; les mises à jour qui surviennent sur un chunk après sa lecture ne sont pas fusionnées pendant la phase de lecture complète et sont lues depuis l'OpLog après l'entrée en phase incrémentielle. Par exemple, une mise à jour de chunk5 survenant pendant le snapshot de chunk5 est reflétée directement dans le snapshot de chunk5 ; si chunk5 est mis à jour après que le lecteur est passé à chunk80, la mise à jour est appliquée ultérieurement depuis l'OpLog lors de la phase incrémentielle.

Important

Lorsque cette option est activée, les modifications survenant pendant ou après l'analyse d'un chunk sont toujours livrées depuis l'OpLog en phase incrémentielle et peuvent être dupliquées. Seule la sémantique at-least-once est garantie. N'activez cette option que si le sink en aval prend en charge les écritures idempotentes par clé primaire.

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

Non

BOOLEAN

false

Lit d'abord les chunks non bornés. Réduit le risque de manque de mémoire pour les collections fréquemment mises à jour.

batch.size

Non

INT

1024

Taille du lot du curseur.

poll.max.batch.size

Non

INT

1024

Nombre maximal d'entrées par requête de récupération du Change Stream.

poll.await.time.ms

Non

INT

1000

Temps d'attente minimal entre les requêtes de récupération du Change Stream, en millisecondes.

heartbeat.interval.ms

Non

INT

0

Intervalle de battement de cœur en millisecondes. Configurez cette option pour les collections rarement mises à jour. Définir sur 0 désactive les battements de cœur.

scan.incremental.snapshot.chunk.size.mb

Non

INT

64

Taille des chunks lors du snapshot, en Mo.

scan.incremental.snapshot.chunk.samples

Non

INT

20

Nombre d'échantillons utilisés pour estimer la taille de la collection lors du snapshot.

scan.full-changelog

Non

BOOLEAN

false

Génère des événements de journal de modifications complets à l'aide des enregistrements preimage et postimage. Nécessite MongoDB 6.0 ou ultérieur avec preimage et postimage activés.

scan.cursor.no-timeout

Non

BOOLEAN

false

Désactive le délai d'expiration du curseur. Par défaut, MongoDB ferme les curseurs inactifs après 10 minutes.

scan.ignore-delete.enabled

Non

BOOLEAN

false

Ignore les événements de suppression provenant de MongoDB.

scan.flatten.nested-documents.enabled

Non

BOOLEAN

false

Aplati les documents BSON imbriqués. Par exemple, {"doc": {"foo": 1, "bar": "two"}} devient doc.foo INT, doc.bar STRING.

scan.all.primitives.as-string.enabled

Non

BOOLEAN

false

Infère tous les types primitifs comme STRING. Réduit les événements de changement de schéma lorsque les types en amont sont incohérents.

metadata.list

Non

STRING

Liste de champs de métadonnées séparés par des virgules à transmettre en aval. Valeurs prises en charge : ts_ms (horodatage de l'événement OpLog), op_ts (alias de ts_ms ; à utiliser lors de l'écriture de métadonnées vers Kafka JSON).

Mappages des types de données

BSON MongoDB

Flink CDC

Notes

STRING

VARCHAR

INT32

INT

INT64

BIGINT

DECIMAL128

DECIMAL

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

TIMESTAMP

TIMESTAMP

DATETIME

LOCALZONEDTIMESTAMP

BINARY

VARBINARY

DOCUMENT

MAP

Les types de clé et de valeur sont inférés.

ARRAY

ARRAY

Les types d'éléments sont inférés.

OBJECTID

VARCHAR

Représenté sous forme de chaîne hexadécimale.

SYMBOL, REGULAREXPRESSION, JAVASCRIPT, JAVASCRIPTWITHSCOPE

VARCHAR

Représenté sous forme de chaîne.

Colonnes de métadonnées

Flink CDC prend en charge la colonne de métadonnées suivante pour le connecteur MongoDB :

Colonne de métadonnées

Type

Description

ts_ms

BIGINT NOT NULL

Heure de modification du document (horodatage OpLog). Retourne 0 pour les documents issus du snapshot initial.

Utilisez les colonnes de métadonnées génériques du module Transform pour accéder à database_name, collection_name et row_kind.

DataStream API

Important

Pour utiliser l'API DataStream, configurez le connecteur DataStream pour votre job. Consultez Utilisation du connecteur DataStream.

Ajouter la dépendance Maven

Le dépôt central Maven héberge les connecteurs MongoDB VVR.

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

Construire un MongoDBSource

Utilisez MongoDBSource.builder() pour construire une source :

  • Pour activer la lecture incrémentielle de snapshot, utilisez le builder de com.ververica.cdc.connectors.mongodb.source.

  • Sinon, utilisez le builder de com.ververica.cdc.connectors.mongodb.

MongoDBSource.builder()
  .hosts("mongo.example.com:27017")
  .username("mongouser")
  .password("mongopasswd")
  .databaseList("testdb")          // Supports regular expressions; use .* to match all databases
  .collectionList("testcoll")      // Supports regular expressions; use .* to match all collections
  .startupOptions(StartupOptions.initial()) // StartupOptions.latest-offset(), StartupOptions.timestamp()
  .deserializer(new JsonDebeziumDeserializationSchema())
  .build();

Paramètres de MongoDBSource

Paramètre

Description

hosts

Nom d'hôte du serveur MongoDB.

username

Nom d'utilisateur MongoDB. À omettre si l'authentification n'est pas activée.

password

Mot de passe MongoDB. À omettre si l'authentification n'est pas activée.

databaseList

Nom de la base de données à surveiller. Prend en charge les expressions régulières. Utilisez .* pour faire correspondre toutes les bases de données.

collectionList

Nom de la collection à surveiller. Prend en charge les expressions régulières. Utilisez .* pour faire correspondre toutes les collections.

startupOptions

Mode de démarrage. Valeurs valides : StartupOptions.initial(), StartupOptions.latest-offset(), StartupOptions.timestamp().

deserializer

Désérialiseur pour convertir les objets SourceRecord. Valeurs valides : MongoDBConnectorDeserializationSchema (mode upsert, produit des RowData Flink), MongoDBConnectorFullChangelogDeserializationSchema (mode journal de modifications complet, produit des RowData Flink), JsonDebeziumDeserializationSchema (produit des chaînes JSON).

Références