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 :
Snapshot complet : lecture parallèle de tous les documents existants dans les collections ciblées.
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 |
|
|
Lit un snapshot au premier démarrage, puis passe en lecture incrémentielle |
Une copie complète des données existantes est nécessaire |
|
|
Démarre à partir de la position actuelle de l'oplog ; aucune donnée historique |
Seules les modifications futures sont requises |
|
|
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.

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.collectionsetconfig.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.enabledsurtrue.Les bases de données
admin,localetconfig, 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 STRINGet 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-guaranteeaccepte les valeursnoneouat-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é
_idde 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 |
|
|
String |
Oui |
— |
Identifiant du connecteur. Tables source : |
|
|
String |
Non |
— |
URI de connexion MongoDB. Spécifiez soit |
|
|
String |
Non |
— |
Nom d'hôte du serveur MongoDB. Séparez plusieurs hôtes par des virgules ( |
|
|
String |
Non |
|
Protocole de connexion. Valeurs valides : |
|
|
String |
Non |
— |
Nom d'utilisateur MongoDB. Obligatoire lorsque l'authentification est activée. |
|
|
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. |
|
|
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 |
|
|
String |
Non |
— |
Nom de la collection MongoDB. Prend en charge les expressions régulières pour les tables source. Notes :
|
|
|
String |
Non |
— |
Options de connexion supplémentaires sous forme de paires 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 |
Source
|
Option |
Type |
Obligatoire |
Valeur par défaut |
Description |
|
|
String |
Non |
|
Mode de démarrage. Valeurs valides : |
|
|
Long |
Conditionnel |
— |
Horodatage de début en millisecondes depuis l'époque UNIX. Obligatoire lorsque |
|
|
Integer |
Non |
|
Taille maximale de la file d'attente pendant la phase de snapshot initial. Ne prend effet que lorsque |
|
|
Integer |
Non |
|
Taille du lot du curseur. |
|
|
Integer |
Non |
|
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. |
|
|
Integer |
Non |
|
Intervalle entre les requêtes de récupération de données, en millisecondes. |
|
|
Integer |
Non |
|
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 |
|
|
Boolean |
Non |
|
Active la lecture parallèle de snapshot. Fonctionnalité expérimentale. Nécessite MongoDB 4.0 ou ultérieur. |
|
|
Integer |
Non |
|
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. |
|
|
Boolean |
Non |
|
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. |
|
|
Boolean |
Non |
|
Analyse les champs séparés par |
|
|
Boolean |
Non |
|
Analyse tous les types BSON primitifs comme STRING. Pris en charge dans VVR 8.0.5 et versions ultérieures. |
|
|
Boolean |
Non |
|
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. |
|
|
Boolean |
Non |
|
Valeurs valides :
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. |
|
|
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 : |
|
|
Integer |
Non |
— |
Nombre de threads pour la réplication du snapshot. Ne prend effet que lorsque |
|
|
Integer |
Non |
|
Taille de la file d'attente pour le snapshot initial. Ne prend effet que lorsque |
|
|
Integer |
Non |
|
Nombre de lecteurs concurrents pour le Change Stream. Ne prend effet que lorsque |
|
|
Integer |
Non |
|
Taille de la file d'attente des messages pour les abonnements concurrents au change stream. Ne prend effet que lorsque |
Lookup (dimension)
|
Option |
Type |
Obligatoire |
Valeur par défaut |
Description |
|
|
String |
Non |
|
Politique de cache. Valeurs valides : |
|
|
Integer |
Non |
|
Nombre maximal de tentatives en cas d'échec du lookup. |
|
|
Duration |
Non |
|
Intervalle entre les tentatives en cas d'échec du lookup. |
|
|
Duration |
Non |
— |
Durée de vie maximale d'une entrée en cache après le dernier accès. Unités prises en charge : |
|
|
Duration |
Non |
— |
Durée de vie maximale d'une entrée en cache après son écriture. Nécessite |
|
|
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 |
|
|
Boolean |
Non |
|
Met en cache une entrée nulle lorsqu'une clé de lookup ne correspond à aucun enregistrement. Nécessite |
Sink
|
Option |
Type |
Obligatoire |
Valeur par défaut |
Description |
|
|
Integer |
Non |
|
Nombre maximal d'enregistrements écrits par lot. |
|
|
Duration |
Non |
|
Intervalle de vidage du tampon. |
|
|
String |
Non |
|
Sémantique de livraison des écritures. Valeurs valides : |
|
|
Integer |
Non |
|
Nombre maximal de tentatives en cas d'échec d'écriture. |
|
|
Duration |
Non |
|
Intervalle entre les tentatives en cas d'échec d'écriture. |
|
|
Integer |
Non |
— |
Parallélisme personnalisé du sink. |
|
|
String |
Non |
|
Stratégie de gestion des événements -D et -U. Valeurs valides : |
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 |
|
|
STRING NOT NULL |
Base de données contenant le document. |
|
|
STRING NOT NULL |
Collection contenant le document. |
|
|
TIMESTAMP_LTZ(3) NOT NULL |
Heure de modification du document. Retourne |
|
|
STRING NOT NULL |
Type d'événement de modification : |
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 |
|
|
Oui |
STRING |
— |
Connecteur. Définissez sur |
|
|
Non |
STRING |
|
Protocole de connexion. Valeurs valides : |
|
|
Oui |
STRING |
— |
Nom(s) d'hôte du serveur MongoDB. Séparez plusieurs hôtes par des virgules. |
|
|
Non |
STRING |
— |
Nom d'utilisateur MongoDB. |
|
|
Non |
STRING |
— |
Mot de passe MongoDB. |
|
|
Oui |
STRING |
— |
Nom de la base de données MongoDB à capturer. Les expressions régulières sont prises en charge. |
|
|
Oui |
STRING |
— |
Nom de la collection MongoDB à capturer. Les expressions régulières sont prises en charge. Utilisez le namespace pleinement qualifié |
|
|
Non |
STRING |
— |
Options de connexion supplémentaires sous forme de paires |
|
|
Non |
STRING |
|
Stratégie d'inférence de schéma. |
|
|
Non |
INT |
|
Nombre maximal d'enregistrements à échantillonner par collection lors de l'inférence initiale du schéma. |
|
|
Non |
STRING |
|
Mode de démarrage. Valeurs valides : |
|
|
Non |
LONG |
— |
Horodatage de début en millisecondes. Obligatoire lorsque |
|
|
Non |
INT |
|
Taille maximale des chunks de métadonnées. |
|
|
Non |
BOOLEAN |
|
Ferme les lecteurs source inactifs après le basculement en lecture incrémentielle. |
|
|
Non |
BOOLEAN |
|
Valeurs valides :
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. |
|
|
Non |
BOOLEAN |
|
Lit d'abord les chunks non bornés. Réduit le risque de manque de mémoire pour les collections fréquemment mises à jour. |
|
|
Non |
INT |
|
Taille du lot du curseur. |
|
|
Non |
INT |
|
Nombre maximal d'entrées par requête de récupération du Change Stream. |
|
|
Non |
INT |
|
Temps d'attente minimal entre les requêtes de récupération du Change Stream, en millisecondes. |
|
|
Non |
INT |
|
Intervalle de battement de cœur en millisecondes. Configurez cette option pour les collections rarement mises à jour. Définir sur |
|
|
Non |
INT |
|
Taille des chunks lors du snapshot, en Mo. |
|
|
Non |
INT |
|
Nombre d'échantillons utilisés pour estimer la taille de la collection lors du snapshot. |
|
|
Non |
BOOLEAN |
|
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. |
|
|
Non |
BOOLEAN |
|
Désactive le délai d'expiration du curseur. Par défaut, MongoDB ferme les curseurs inactifs après 10 minutes. |
|
|
Non |
BOOLEAN |
|
Ignore les événements de suppression provenant de MongoDB. |
|
|
Non |
BOOLEAN |
|
Aplati les documents BSON imbriqués. Par exemple, |
|
|
Non |
BOOLEAN |
|
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. |
|
|
Non |
STRING |
— |
Liste de champs de métadonnées séparés par des virgules à transmettre en aval. Valeurs prises en charge : |
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 |
|
|
BIGINT NOT NULL |
Heure de modification du document (horodatage OpLog). Retourne |
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
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 |
|
|
Nom d'hôte du serveur MongoDB. |
|
|
Nom d'utilisateur MongoDB. À omettre si l'authentification n'est pas activée. |
|
|
Mot de passe MongoDB. À omettre si l'authentification n'est pas activée. |
|
|
Nom de la base de données à surveiller. Prend en charge les expressions régulières. Utilisez |
|
|
Nom de la collection à surveiller. Prend en charge les expressions régulières. Utilisez |
|
|
Mode de démarrage. Valeurs valides : |
|
|
Désérialiseur pour convertir les objets |
Références
Instruction CREATE TABLE AS (CTAS) (en retrait) — Synchronise les données MongoDB et les changements de schéma vers des tables en aval (VVR 8.0.6 et ultérieur, nécessite preimage et postimage).
Instruction CREATE DATABASE AS (CDAS) — Synchronise une base de données MongoDB entière vers des tables en aval (VVR 8.0.6 et ultérieur, nécessite preimage et postimage).
Métriques — Surveille les performances de la table source.