Une fois le catalogue MongoDB configuré, accédez directement aux collections MongoDB dans vos déploiements Flink sans définir leurs schémas.
Informations générales
Un catalogue MongoDB déduit automatiquement le schéma d'une collection en analysant les documents BSON. Vous récupérez ainsi les informations sur les champs sans déclarer le schéma dans Flink SQL. Un catalogue MongoDB offre les fonctionnalités suivantes :
Aucun enregistrement manuel des tables via une instruction DDL n'est nécessaire, ce qui améliore l'efficacité et la précision du développement.
Les tables issues d'un catalogue MongoDB servent directement de tables source, de tables de dimension et de tables de résultat dans les déploiements Flink SQL.
Dans Ververica Runtime (VVR) 8.0.6 et versions ultérieures, utilisez un catalogue MongoDB avec une instruction CREATE TABLE AS (CTAS) ou une instruction CREATE DATABASE AS (CDAS) pour synchroniser les modifications de schéma.
Cette rubrique décrit comment gérer un catalogue MongoDB :
Limites
Seules les versions Ververica Runtime (VVR) 8.0.5 et ultérieures prennent en charge les catalogues MongoDB.
Vous ne pouvez pas modifier un catalogue MongoDB existant à l'aide d'une instruction DDL.
Vous pouvez uniquement interroger les tables. La création, la modification ou la suppression de bases de données et de tables est impossible.
Créer un catalogue MongoDB
-
Dans l'éditeur de l'onglet Scripts, saisissez l'instruction de configuration du catalogue MongoDB.
CREATE CATALOG <yourcatalogname> WITH( 'type'='mongodb', 'default-database'='<dbName>', 'hosts'='<hosts>', 'scheme'='<scheme>', 'username'='<username>', 'password'='<password>', 'connection.options'='<connectionOptions>', 'max.fetch.records'='100', 'scan.flatten-nested-columns.enable'='<flattenNestedColumns>', 'scan.primitive-as-string'='<primitiveAsString>' );Paramètre
Type
Description
Obligatoire
Remarques
yourcatalogname
String
Nom du catalogue MongoDB.
Oui
Spécifiez un nom personnalisé en anglais.
ImportantAprès avoir remplacé le paramètre par le nom de votre catalogue, supprimez les chevrons (<>). Sinon, la vérification de la syntaxe échoue.
type
String
Type du catalogue.
Oui
Définissez la valeur sur
mongodb.hosts
String
Nom d'hôte du serveur MongoDB.
Oui
Séparez plusieurs noms d'hôte par des virgules (
,).default-database
String
Nom de la base de données MongoDB par défaut.
Oui
Aucun.
scheme
String
Protocole de connexion utilisé pour MongoDB.
Non
Valeurs valides :
-
mongodb(par défaut) : connexion via le protocole MongoDB standard. -
mongodb+srv: connexion via le protocole d'enregistrement DNS SRV.
username
String
Nom d'utilisateur utilisé pour se connecter à MongoDB.
Non
Ce paramètre est requis si l'authentification est activée.
password
String
Mot de passe utilisé pour se connecter à MongoDB.
Non
Ce paramètre est requis si l'authentification est activée.
RemarquePour éviter d'exposer le mot de passe, nous vous recommandons d'utiliser des variables. Pour plus d'informations, consultez Gérer les variables.
connection.options
String
Paramètres de connexion supplémentaires pour le client MongoDB.
Non
Options supplémentaires au format
clé=valeur, séparées par des esperluettes (&). Exemple :connectTimeoutMS=12000&socketTimeoutMS=13000.max.fetch.records
Int
Nombre maximal de documents à extraire pour la déduction du schéma à partir des documents BSON.
Non
Valeur par défaut : 100.
scan.flatten-nested-columns.enabled
Boolean
Indique s'il faut aplatir récursivement les documents imbriqués dans BSON.
Non
Valeurs valides :
-
true: Aplatit récursivement les colonnes imbriquées. Pour une colonne aplatie, Flink utilise le chemin d'accès comme nom de colonne. Par exemple, pour la colonnecoldans{"nested":{"col":true}}, le nom de la colonne aplatie estnested.col. -
false(par défaut) : traite les documents BSON imbriqués comme STRING.
ImportantCe paramètre est pris en charge uniquement lorsqu'une table issue du catalogue MongoDB sert de table source dans un déploiement Flink SQL.
scan.primitive-as-string
Boolean
Indique s'il faut déduire tous les types primitifs comme STRING lors de l'analyse des documents BSON.
Non
Valeurs valides :
-
true: déduit tous les types primitifs comme STRING. -
false(par défaut) : déduit les types selon les règles par défaut. Pour plus d'informations, consultez Détails des tables issus d'un catalogue MongoDB.
-
-
Sélectionnez l'instruction et cliquez sur Run dans la gouttière.
CREATE CATALOG MongoDBCatalog WITH( 'type'='mongodb', 'default-database'='<dbName>', 'hosts'='<hosts>', 'scheme'='<scheme>', 'username'='<username>', 'password'='<password>', 'connection.options'='<connectionOptions>', 'max.fetch.records'='100', 'scan.flatten-nested-columns.enable'='<flattenNestedColumns>', 'scan.primitive-as-string'='<primitiveAsString>' ); Dans le volet Catalogs sur la gauche, vérifiez que le nouveau catalogue apparaît.
Afficher un catalogue MongoDB
-
Dans l'éditeur de l'onglet Scripts, saisissez la commande suivante.
DESCRIBE `${catalog_name}`.`${db_name}`.`${collection_name}`;Paramètre
Description
${catalog_name}
Nom du catalogue MongoDB.
${db_name}
Nom de la base de données MongoDB.
${collection_name}
Nom de la collection MongoDB.
Sélectionnez l'instruction et cliquez sur Run dans la gouttière.
Utiliser un catalogue MongoDB
-
En tant que table source pour lire les données depuis MongoDB.
INSERT INTO ${other_sink_table} SELECT... FROM `${mongodb_catalog}`.`${db_name}`.`${collection_name}` /*+OPTIONS('scan.incremental.snapshot.enabled'='true')*/;RemarquePour spécifier d'autres options WITH pour une table issue d'un catalogue MongoDB, utilisez une indication SQL. Par exemple, l'instruction SQL précédente utilise une indication SQL pour activer l'analyse parallèle pour le snapshot initial. Pour plus d'informations sur les autres options WITH, consultez MongoDB.
-
En tant que table source, utilisez une instruction CREATE TABLE AS (CTAS) ou une instruction CREATE DATABASE AS (CDAS) pour synchroniser les données de MongoDB vers une table de destination.
ImportantPour utiliser une instruction CTAS ou CDAS afin de synchroniser les données de MongoDB vers une table de destination, les conditions suivantes s'appliquent :
La version de VVR doit être 8.0.6 ou ultérieure, et la version de MongoDB doit être 6,0 ou ultérieure.
Les paramètres
scan.incremental.snapshot.enabledetscan.full-changelogdans l'indication SQL doivent être définis surtrue.La fonctionnalité d'images pré et post doit être activée pour la base de données MongoDB. Pour plus d'informations, consultez Document Preimages.
-
Synchronisez une seule table en temps réel.
CREATE TABLE IF NOT EXISTS `${target_table_name}` WITH(...) AS TABLE `${mongodb_catalog}`.`${db_name}`.`${collection_name}` /*+ OPTIONS('scan.incremental.snapshot.enabled'='true', 'scan.full-changelog'='true') */; -
Synchronisez plusieurs tables dans un seul déploiement.
BEGIN STATEMENT SET; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table0` AS TABLE `mongodb-catalog`.`database`.`collection0` /*+ OPTIONS('scan.incremental.snapshot.enabled'='true', 'scan.full-changelog'='true') */; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table1` AS TABLE `mongodb-catalog`.`database`.`collection1` /*+ OPTIONS('scan.incremental.snapshot.enabled'='true', 'scan.full-changelog'='true') */; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table2` AS TABLE `mongodb-catalog`.`database`.`collection2` /*+ OPTIONS('scan.incremental.snapshot.enabled'='true', 'scan.full-changelog'='true') */; END;Vous pouvez utiliser un catalogue MongoDB pour synchroniser plusieurs collections MongoDB dans un seul déploiement, à condition que les conditions suivantes soient remplies :
Toutes les tables doivent avoir les mêmes configurations MongoDB, y compris
hosts,scheme,username,passwordetconnection.options.Toutes les tables doivent avoir la même configuration
scan.startup.mode.
-
Synchronisez une base de données entière.
CREATE DATABASE IF NOT EXISTS `some_catalog`.`some_database` AS DATABASE `mongodb-catalog`.`database` /*+ OPTIONS('scan.incremental.snapshot.enabled'='true', 'scan.full-changelog'='true') */;
-
Lisez les données depuis une table de dimension MongoDB.
INSERT INTO ${other_sink_table} SELECT ... FROM ${other_source_table} AS e JOIN `${mongodb_catalog}`.`${db_name}`.`${table_name}` FOR SYSTEM_TIME AS OF e.proctime AS w ON e.id = w.id; -
Écrivez les données de résultat dans une table MongoDB.
INSERT INTO `${mongodb_catalog}`.`${db_name}`.`${table_name}` SELECT ... FROM ${other_source_table}
Une fois l'instruction exécutée avec succès, affichez les détails de la table dans les résultats d'exécution.
Le schéma de la table comprend les champs suivants : _id (STRING, clé primaire, NOT NULL), name (STRING, nullable), age (INT, nullable) et addr (STRING, nullable).
Supprimer un catalogue MongoDB
La suppression d'un catalogue MongoDB n'affecte pas les déploiements en cours d'exécution. Toutefois, les tentatives de démarrage ou de redémarrage des déploiements utilisant le catalogue supprimé échoueront car leurs tables seront inaccessibles.
-
Dans l'éditeur de l'onglet Scripts, saisissez la commande suivante.
DROP CATALOG ${catalog_name};Dans cette commande, ${catalog_name} spécifie le nom du catalogue MongoDB à supprimer.
Sélectionnez l'instruction, cliquez dessus avec le bouton droit de la souris, puis choisissez Run.
Dans le volet Catalogs sur la gauche, vérifiez que le catalogue n'est plus répertorié.
Détails des tables issus d'un catalogue MongoDB
Pour simplifier l'utilisation, un catalogue MongoDB ajoute automatiquement des configurations par défaut et une clé primaire aux tables déduites. Pour déduire le schéma d'une collection, le catalogue extrait jusqu'à max.fetch.records documents, analyse chacun d'eux et fusionne les résultats dans un schéma final. Le schéma comprend les parties suivantes :
-
Colonnes physiques déduites
Un catalogue MongoDB déduit les colonnes physiques à partir des documents BSON.
-
Contrainte de clé primaire par défaut
Pour les tables issues d'un catalogue MongoDB, la colonne
_idsert de clé primaire par défaut pour éviter les données en double.
Après avoir extrait un ensemble de documents BSON, le catalogue les analyse individuellement et fusionne les colonnes physiques résultantes selon les règles suivantes pour former le schéma final de la collection :
Si une colonne physique analysée contient un champ absent du schéma de résultat, le catalogue MongoDB ajoute automatiquement le champ au schéma de résultat.
-
Si plusieurs colonnes partagent le même nom, le conflit est résolu comme suit :
Si les types de données sont identiques mais que la précision diffère, le type ayant la précision la plus élevée est utilisé.
-
Si les types de données sont différents, l'ancêtre commun le plus bas dans l'arborescence hiérarchique des types illustrée dans la figure suivante est utilisé comme type pour la colonne. Toutefois, pour préserver la précision lors de la fusion des types Decimal et Float, le type de résultat est Double.
Lors de la déduction du schéma, les types de données BSON correspondent aux types de données Flink SQL comme suit :
|
Type BSON |
Type Flink SQL |
|
Boolean |
BOOLEAN |
|
Int32 |
INT |
|
Int64 |
BIGINT |
|
Binary |
BYTES |
|
Double |
DOUBLE |
|
Decimal128 |
DECIMAL |
|
String |
STRING |
|
ObjectId |
STRING |
|
DateTime |
TIMESTAMP_LTZ(3) |
|
Timestamp |
TIMESTAMP_LTZ(0) |
|
Array |
STRING |
|
Document |
STRING |
Documentation connexe
Pour plus d'informations sur le connecteur MongoDB, consultez MongoDB.
Si les catalogues intégrés ne répondent pas à vos besoins métier, utilisez des catalogues personnalisés. Pour plus d'informations, consultez Gérer les catalogues personnalisés.