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 :
Flink échantillonne jusqu'à 100 messages du topic
Il déduit le schéma à partir de la structure JSON
-
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.
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.
RemarqueDans 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
-
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' ); -
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é.
ImportantSupprimez 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.
RemarqueLa 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.
RemarqueLa 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.RemarqueSi 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.
RemarqueSeules 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.
RemarqueSeules 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.
RemarqueSeules 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.
RemarqueSeules 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.
RemarqueSeules 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.
RemarqueSeules les versions VVR 6.0.2 ou ultérieures prennent en charge ce paramètre.
-
-
Sélectionnez l'instruction CREATE CATALOG, puis cliquez sur Run.

Vérifiez que le catalogue apparaît dans la zone Catalogs située à gauche.
Afficher un catalogue Kafka JSON
-
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.
-
Exécutez l'instruction pour afficher le schéma.

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') */;
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 :

Dans cet exemple, les deux premières tables répondent à toutes les exigences, tandis que les deux dernières ne les respectent pas.
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
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 ».
-
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.
Sélectionnez l'instruction DROP CATALOG, effectuez un clic droit, puis sélectionnez Run.
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.
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.
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.

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

-
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,
partitionetoffsetservent de clé primaire pour garantir l'unicité des données.
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 :
kafkaouupsert-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.