Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Gérer les catalogues MongoDB

Dernière mise à jour :Aug 19, 2026

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

  1. 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.

    Important

    Aprè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.

    Remarque

    Pour é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 colonne col dans {"nested":{"col":true}}, le nom de la colonne aplatie est nested.col.

    • false (par défaut) : traite les documents BSON imbriqués comme STRING.

    Important

    Ce 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 :

  2. 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>'
    );
  3. Dans le volet Catalogs sur la gauche, vérifiez que le nouveau catalogue apparaît.

Afficher un catalogue MongoDB

  1. 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.

  2. 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')*/;
    Remarque

    Pour 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.

    Important

    Pour 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.enabled et scan.full-changelog dans l'indication SQL doivent être définis sur true.

    • 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, password et connection.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

Avertissement

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.

  1. 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.

  2. Sélectionnez l'instruction, cliquez dessus avec le bouton droit de la souris, puis choisissez Run.

  3. 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 _id sert 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.

      image

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.