Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Gérer les catalogues Kafka JSON

Dernière mise à jour :Aug 09, 2026

Déduisez automatiquement les schémas de table à partir des messages JSON Kafka et interrogez les topics sans écrire d'instructions DDL manuelles.

Qu'est-ce qu'un catalogue Kafka JSON ?

Un catalogue Kafka JSON déduit automatiquement la structure des tables à partir des messages Kafka au format JSON. Cela vous permet d'interroger les topics avec SQL, sans avoir à rédiger d'instructions DDL.

Principaux avantages :

  • Développement plus rapide : Interrogez directement les topics Kafka sans déclarer de schémas

  • Moins d'erreurs : Les noms de table correspondent automatiquement aux noms de topic

  • Évolution du schéma : Utilisez-le avec CTAS pour synchroniser les données lorsque les schémas évoluent

Exemple : Au lieu d'écrire cette instruction DDL :

CREATE TABLE orders (
  order_id STRING,
  product_name STRING,
  quantity INT,
  price DECIMAL(10,2)
) WITH (
  'connector' = 'kafka',
  'topic' = 'orders',
  'properties.bootstrap.servers' = '...',
  'format' = 'json'
);

Vous pouvez interroger directement le topic :

SELECT * FROM kafka_catalog.kafka.orders;

Fonctionnement

Lorsque vous interrogez une table de catalogue :

  1. Flink échantillonne jusqu'à 100 messages du topic

  2. Il déduit le schéma à partir de la structure JSON

  3. Il crée une table comportant :

    • Les champs issus de votre message JSON (clé et valeur)

    • Des colonnes de métadonnées (partition, offset, horodatage)

    • Une clé primaire (partition, offset)

Ce processus est automatique : aucune instruction DDL n'est requise.

Important

Pour les topics contenant des schémas mixtes, le catalogue fusionne tous les champs en un seul schéma. Détails sur l'inférence de schéma.

Dans cette rubrique :

Limites

  • Seuls les messages au format JSON sont pris en charge.

  • Requiert VVR 6.0.2 ou une version ultérieure.

  • La modification d'un catalogue Kafka JSON n'est pas prise en charge.

  • Les tables de catalogue sont en lecture seule.

    Remarque

    Dans les scénarios CREATE DATABASE AS (CDAS) ou CREATE TABLE AS (CTAS) utilisant un catalogue Kafka JSON, les topics peuvent être créés automatiquement.

  • Les catalogues Kafka JSON ne peuvent ni lire ni écrire dans des clusters Kafka disposant de l'authentification SSL ou SASL activée.

  • Les tables fournies par les catalogues Kafka JSON peuvent être utilisées directement comme tables source dans les jobs Flink SQL. Elles ne peuvent pas servir de tables sink ni de tables de dimension pour les recherches (lookup).

  • ApsaraMQ for Kafka ne permet pas actuellement de supprimer des groupes de consommateurs via la même opération API qu'Apache Kafka. Lors de la création d'un catalogue Kafka JSON, vous devez configurer les paramètres aliyun.kafka.instanceId, aliyun.kafka.accessKeyId, aliyun.kafka.accessKeySecret, aliyun.kafka.endpoint et aliyun.kafka.regionId afin de supprimer automatiquement les groupes de consommateurs. Comparaison entre ApsaraMQ for Kafka et Apache Kafka.

Notes d'utilisation

Cohérence du schéma : Le catalogue échantillonne les messages pour déduire les schémas. Si les messages présentent des formats différents, le catalogue fusionne tous les champs en un seul schéma.

Impact des modifications de schéma : Si le format des messages change, les redémarrages de job peuvent échouer car le plan d'exécution utilise l'ancien schéma.

Solution : Corrigez le schéma avec CREATE TEMPORARY TABLE. Exemple :

-- Define a fixed schema based on the catalog table
CREATE TEMPORARY TABLE orders (
  value_order_id STRING,
  value_product_name STRING,
  value_quantity INT,
  value_price DECIMAL(10,2)
) LIKE `kafka_catalog`.`kafka`.`orders`;

Cette approche verrouille le schéma et empêche les échecs lors du redémarrage.

Créer un catalogue Kafka JSON

  1. Dans l'éditeur SQL de la page Scripts, saisissez l'instruction permettant de créer un catalogue Kafka JSON.

    • Cluster Kafka autogéré ou Cluster Kafka EMR on ECS

      CREATE CATALOG <YourCatalogName> WITH(
       'type'='kafka',  -- Required
       'properties.bootstrap.servers'='<brokers>',  -- Required
       'format'='json',  -- Required
       'default-database'='<dbName>',
       'key.fields-prefix'='<keyPrefix>',
       'value.fields-prefix'='<valuePrefix>',
       'timestamp-format.standard'='<timestampFormat>',
       'infer-schema.flatten-nested-columns.enable'='<flattenNestedColumns>',
       'infer-schema.primitive-as-string'='<primitiveAsString>',
       'infer-schema.parse-key-error.field-name'='<parseKeyErrorFieldName>',
       'infer-schema.compacted-topic-as-upsert-table'='true',
       'max.fetch.records'='100'
      );
    • ApsaraMQ for Kafka

      CREATE CATALOG <YourCatalogName> WITH(
       'type'='kafka',  -- Required
       'properties.bootstrap.servers'='<brokers>',  -- Required
       'format'='json',  -- Required
       'default-database'='<dbName>',
       'key.fields-prefix'='<keyPrefix>',
       'value.fields-prefix'='<valuePrefix>',
       'timestamp-format.standard'='<timestampFormat>',
       'infer-schema.flatten-nested-columns.enable'='<flattenNestedColumns>',
       'infer-schema.primitive-as-string'='<primitiveAsString>',
       'infer-schema.parse-key-error.field-name'='<parseKeyErrorFieldName>',
       'infer-schema.compacted-topic-as-upsert-table'='true',
       'max.fetch.records'='100',
       'aliyun.kafka.accessKeyId'='<aliyunAccessKeyId>',  -- Required
       'aliyun.kafka.accessKeySecret'='<aliyunAccessKeySecret>',  -- Required
       'aliyun.kafka.instanceId'='<aliyunKafkaInstanceId>',  -- Required
       'aliyun.kafka.endpoint'='<aliyunKafkaEndpoint>',  -- Required
       'aliyun.kafka.regionId'='<aliyunKafkaRegionId>'  -- Required
      );

    Paramètre

    Type

    Description

    Obligatoire

    Remarques

    YourCatalogName

    String

    Nom du catalogue.

    Oui

    Saisissez un nom personnalisé.

    Important

    Supprimez les chevrons (<>) lorsque vous remplacez l'espace réservé.

    type

    String

    Type de catalogue.

    Oui

    Doit être kafka.

    properties.bootstrap.servers

    String

    Adresses des brokers Kafka.

    Oui

    Format : host1:port1,host2:port2,host3:port3.

    Séparez les différentes adresses par des virgules (,).

    format

    String

    Format des messages Kafka.

    Oui

    Doit être json.

    default-database

    String

    Nom du cluster.

    Non

    Par défaut : kafka. Définit db_name dans le nom en trois parties catalog_name.db_name.table_name. Kafka ne possédant pas de bases de données, utilisez n'importe quelle chaîne de caractères.

    key.fields-prefix

    String

    Préfixe pour les champs de clé de message. Évite les conflits de nommage.

    Non

    Par défaut : key_. Exemple : Le champ a devient key_a.

    Remarque

    La valeur du paramètre key.fields-prefix ne peut pas être un préfixe de la valeur du paramètre value.fields-prefix. Par exemple, si vous définissez value.fields-prefix sur test1_value_, vous ne pouvez pas définir key.fields-prefix sur test1_.

    value.fields-prefix

    String

    Préfixe pour les champs de valeur de message. Évite les conflits de nommage.

    Non

    Par défaut : value_. Exemple : Le champ b devient value_b.

    Remarque

    La valeur du paramètre value.fields-prefix ne peut pas être un préfixe de la valeur du paramètre key.fields-prefix. Par exemple, si vous définissez key.fields-prefix sur test2_value_, vous ne pouvez pas définir value.fields-prefix sur test2_.

    timestamp-format.standard

    String

    Format d'horodatage pour les messages JSON. Flink tente d'abord le format configuré, puis revient à d'autres formats en cas d'échec.

    Non

    Valeurs valides :

    • SQL (par défaut)

    • ISO-8601

    infer-schema.flatten-nested-columns.enable

    Boolean

    Développe récursivement les colonnes imbriquées dans les valeurs de message.

    Non

    Valeurs valides :

    • true : Développe les colonnes imbriquées.

      Les noms des colonnes développées utilisent le chemin comme nom. Exemple : col dans {"nested": {"col": true}} devient nested.col.

      Remarque

      Si vous définissez ce paramètre sur true, utilisez-le conjointement avec l'instruction CREATE TABLE AS (CTAS). Les autres instructions DML ne prennent pas en charge le développement automatique des colonnes imbriquées.

    • false (par défaut) : Traite les types imbriqués comme String.

    infer-schema.primitive-as-string

    Boolean

    Déduit tous les types primitifs comme String.

    Non

    Valeurs valides :

    • true : Déduit tous les types primitifs comme String.

    • false (par défaut) : Déduit les types selon les règles de base.

    infer-schema.parse-key-error.field-name

    String

    Si la clé du message n'est pas vide mais ne peut pas être analysée, un champ VARBINARY est ajouté. Le nom du champ combine le préfixe key.fields-prefix avec la valeur de ce paramètre.

    Non

    Par défaut : col. Exemple : Si la valeur est analysée en value_name et que la clé échoue à l'analyse, le schéma contient key_col et value_name.

    infer-schema.compacted-topic-as-upsert-table

    Boolean

    Traite la table comme une table Upsert Kafka lorsque la politique de nettoyage du topic est compact et que la clé du message n'est pas vide.

    Non

    Par défaut : true. Activez cette option lors de l'utilisation de CTAS ou CDAS pour synchroniser les données vers ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

    max.fetch.records

    Int

    Nombre maximal de messages échantillonnés pour l'inférence de schéma.

    Non

    Par défaut : 100.

    aliyun.kafka.accessKeyId

    String

    ID AccessKey de votre compte Alibaba Cloud. Pour plus d'informations, consultez Créer une paire de clés AccessKey.

    Non

    Requis pour les clusters ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

    aliyun.kafka.accessKeySecret

    String

    Secret AccessKey de votre compte Alibaba Cloud. Pour plus d'informations, consultez Créer une paire de clés AccessKey.

    Non

    Requis pour les clusters ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

    aliyun.kafka.instanceId

    String

    ID d'instance ApsaraMQ for Kafka. Vous le trouverez sur la page des détails de l'instance dans la console ApsaraMQ for Kafka.

    Non

    Requis pour les clusters ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

    aliyun.kafka.endpoint

    String

    Point de terminaison API ApsaraMQ for Kafka. Pour plus d'informations, consultez Points de terminaison.

    Non

    Requis pour les clusters ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

    aliyun.kafka.regionId

    String

    ID de région de l'instance Kafka. Pour plus d'informations, consultez Points de terminaison.

    Non

    Requis pour les clusters ApsaraMQ for Kafka.

    Remarque

    Seules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.

  2. Sélectionnez l'instruction CREATE CATALOG, puis cliquez sur Run.

    image.png

  3. Vérifiez que le catalogue apparaît dans la zone Catalogs située à gauche.

Afficher un catalogue Kafka JSON

  1. Dans l'éditeur SQL de la page Scripts, saisissez l'instruction suivante.

    DESCRIBE `${catalog_name}`.`${db_name}`.`${topic_name}`;

    Paramètre

    Description

    ${catalog_name}

    Le nom du catalogue.

    ${db_name}

    Le nom du cluster.

    ${topic_name}

    Le nom du topic.

  2. Exécutez l'instruction pour afficher le schéma.

    Table information

Utiliser un catalogue Kafka JSON

Après avoir créé un catalogue Kafka JSON, référencez ses topics dans vos jobs Flink SQL.

Utilisation en tant que table source

Scénario : Extraire des données Kafka et les écrire dans un autre système.

Mode d'emploi : Interrogez directement la table du catalogue :

-- Insert data from Kafka topic to target table
INSERT INTO ${other_sink_table}
SELECT order_id, product_name, quantity * price AS total_amount
FROM `${kafka_catalog}`.`${db_name}`.`${topic_name}`
/*+ OPTIONS('scan.startup.mode'='earliest-offset') */;
Remarque

Utilisez les Indices SQL pour spécifier les options de table, telles que scan.startup.mode. Toutes les options disponibles sont répertoriées dans les Paramètres de la table source Kafka.

Utilisation avec CTAS pour synchroniser les données

Scénario : Synchroniser des topics Kafka entiers sans rédiger d'instructions DDL.

Fonctionnement : L'instruction CREATE TABLE AS (CTAS) (en cours de retrait) crée une table cible dotée du même schéma que la source. Cette approche est utile pour :

  • Synchroniser sans définition manuelle du schéma

  • Gérer automatiquement les modifications de schéma

Synchroniser un seul topic :

-- Create target table with inferred schema
CREATE TABLE IF NOT EXISTS `${target_table_name}`
WITH (
  'connector' = 'hologres',
  'dbname' = 'my_database',
  'tablename' = 'orders'
)
AS TABLE `${kafka_catalog}`.`${db_name}`.`${topic_name}`
/*+ OPTIONS('scan.startup.mode'='earliest-offset') */;

Synchroniser plusieurs topics dans un seul job :

BEGIN STATEMENT SET;

CREATE TABLE IF NOT EXISTS `target_catalog`.`target_db`.`orders`
AS TABLE `kafka_catalog`.`kafka`.`orders`
/*+ OPTIONS('scan.startup.mode'='earliest-offset') */;

CREATE TABLE IF NOT EXISTS `target_catalog`.`target_db`.`products`
AS TABLE `kafka_catalog`.`kafka`.`products`
/*+ OPTIONS('scan.startup.mode'='earliest-offset') */;

CREATE TABLE IF NOT EXISTS `target_catalog`.`target_db`.`customers`
AS TABLE `kafka_catalog`.`kafka`.`customers`
/*+ OPTIONS('scan.startup.mode'='earliest-offset') */;

END;

Conditions requises pour la synchronisation de plusieurs topics :

Lors de la synchronisation de plusieurs topics Kafka dans le même job, assurez-vous que :

  • Aucune des tables n'utilise le paramètre topic-pattern

  • Toutes les tables partagent la même configuration Kafka (mêmes propriétés properties.bootstrap.servers, properties.group.id, etc.)

  • Toutes les tables utilisent le même mode scan.startup.mode (group-offsets, latest-offset ou earliest-offset)

Exemple : L'image suivante illustre les configurations qui respectent ces exigences :

Configuration requirements for syncing multiple topics

Dans cet exemple, les deux premières tables répondent à toutes les exigences, tandis que les deux dernières ne les respectent pas.

Remarque

Pour des exemples complets de bout en bout, consultez le Guide de démarrage rapide pour l'entreposage de journaux en temps réel.

Supprimer un catalogue Kafka JSON

Avertissement

La suppression d'un catalogue n'affecte pas les jobs en cours d'exécution. Toutefois, le déploiement ou le redémarrage de jobs utilisant ce catalogue échouera avec des erreurs « table not found ».

  1. Dans l'éditeur SQL de la page Scripts, saisissez l'instruction suivante.

    DROP CATALOG ${catalog_name};

    Remplacez ${catalog_name} par le nom de votre catalogue.

  2. Sélectionnez l'instruction DROP CATALOG, effectuez un clic droit, puis sélectionnez Run.

  3. Vérifiez que le catalogue n'apparaît plus dans la zone Catalogs.

Référence : Détails sur l'inférence de schéma

Cette section fournit des détails techniques sur la manière dont les catalogues Kafka JSON déduisent les schémas de table.

Remarque

Vous pouvez ignorer cette section si vous souhaitez simplement utiliser le catalogue. Lisez-la lorsque :

  • Vous rencontrez des problèmes liés au schéma

  • Vous souhaitez comprendre comment le catalogue gère les schémas incohérents

  • Vous cherchez à optimiser les performances de l'inférence de schéma

Processus d'inférence de schéma

Lorsque vous interrogez un topic Kafka, Flink échantillonne les messages (jusqu'à max.fetch.records, 100 par défaut) et fusionne leurs schémas.

Processus détaillé : Flink analyse chaque message et fusionne les schémas.

Important
  • L'inférence de schéma crée un groupe de consommateurs (avec un préfixe spécifique au catalogue) pour consommer les données du topic.

  • Pour ApsaraMQ for Kafka, utilisez VVR 6.0.7 ou une version ultérieure. Les versions antérieures ne suppriment pas automatiquement les groupes de consommateurs, ce qui provoque des alertes d'empilement de messages.

Le schéma comprend les colonnes physiques déduites, les colonnes de métadonnées et les contraintes de clé primaire :

  • Colonnes physiques déduites

    Flink déduit les colonnes physiques à partir de la clé et de la valeur du message, en ajoutant aux noms de colonne le préfixe configuré.

    Si la clé n'est pas vide mais ne peut pas être analysée, Flink crée une colonne VARBINARY. Le nom de la colonne combine key.fields-prefix avec la valeur de infer-schema.parse-key-error.field-name.

    Règles de fusion de schéma :

    • Les nouveaux champs sont ajoutés au schéma final.

    • Pour les champs portant le même nom :

      • Mêmes types, précisions différentes : Utiliser la précision la plus élevée.

      • Types différents : Trouver le nœud parent le plus petit dans l'arbre des types (voir figure). Decimal + Float fusionnent en Double pour préserver la précision.Schema merge

    Exemple : Pour un topic contenant ces trois messages, le catalogue produit ce schéma :Schema

  • Colonnes de métadonnées par défaut

    Flink ajoute par défaut trois colonnes de métadonnées : partition, offset et timestamp.

    Nom de la métadonnée

    Nom de la colonne

    Type

    Description

    partition

    partition

    INT NOT NULL

    Numéro de partition.

    offset

    offset

    BIGINT NOT NULL

    Offset du message.

    timestamp

    timestamp

    TIMESTAMP_LTZ(3) NOT NULL

    Horodatage du message.

  • Contrainte PRIMARY KEY par défaut

    Lors de la lecture depuis Kafka, partition et offset servent de clé primaire pour garantir l'unicité des données.

Remarque

Si le schéma déduit ne répond pas à vos besoins, utilisez CREATE TEMPORARY TABLE ... LIKE pour spécifier explicitement un schéma personnalisé. Exemple : Si le JSON contient un champ ts au format '2023-01-01 12:00:01', le catalogue le déduit comme TIMESTAMP. Pour l'utiliser comme STRING, déclarez la table comme indiqué ci-dessous. Notez le préfixe value_ pour les champs de valeur de message :

CREATE TEMPORARY TABLE tempTable (
    value_name STRING,
    value_ts STRING
) LIKE `kafkaJsonCatalog`.`kafka`.`testTopic`;
  • Paramètres de table par défaut

    Paramètre

    Description

    Remarques

    connector

    Type de connecteur.

    Valeur : kafka ou upsert-kafka.

    topic

    Nom du topic.

    Identique au nom de la table.

    properties.bootstrap.servers

    Adresses des brokers Kafka.

    Identique à properties.bootstrap.servers du catalogue.

    value.format

    Format de sérialisation des valeurs de message.

    Toujours JSON.

    value.fields-prefix

    Préfixe pour les champs de valeur de message afin d'éviter les conflits de nommage.

    Identique à value.fields-prefix du catalogue.

    value.json.infer-schema.flatten-nested-columns.enable

    Développe récursivement les colonnes imbriquées dans les valeurs de message.

    Identique à infer-schema.flatten-nested-columns.enable du catalogue.

    value.json.infer-schema.primitive-as-string

    Déduit tous les types primitifs comme String pour les valeurs de message.

    Identique à infer-schema.primitive-as-string du catalogue.

    value.fields-include

    Politique de gestion des champs de clé dans les valeurs de message.

    Doit être EXCEPT_KEY, ce qui signifie que les valeurs de message excluent les champs de clé.

    Vous devez configurer ce paramètre si la clé du message n'est pas vide. Ne configurez pas ce paramètre si la clé du message est vide.

    key.format

    Format utilisé par le connecteur Flink Kafka pour sérialiser ou désérialiser la clé du message Kafka.

    Doit être json ou raw.

    Requis lorsque la clé du message n'est pas vide.

    Si la clé du message n'est pas vide mais ne peut pas être analysée, définissez ce paramètre sur raw. Si l'analyse réussit, définissez ce paramètre sur json.

    key.fields-prefix

    Préfixe pour les champs de clé de message afin d'éviter les conflits de nommage.

    Identique à key.fields-prefix du catalogue.

    Requis lorsque la clé du message n'est pas vide.

    key.fields

    Champs stockant les données de clé de message analysées.

    Rempli automatiquement.

    Requis lorsque la clé du message n'est pas vide et que la table n'est pas une table Upsert Kafka.