Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Kafka

Dernière mise à jour :Aug 27, 2026

Général

Option

Description

Type

Obligatoire

Valeur par défaut

Remarques

connector

Type de connecteur.

String

Oui

La valeur doit être kafka.

properties.bootstrap.servers

Liste des adresses des courtiers Kafka.

String

Oui

Format : host1:port1,host2:port2,.... Séparez les adresses par des virgules (,).

properties.*

Propriétés supplémentaires pour le client Kafka.

String

Non

Les clés de propriété doivent correspondre aux options valides définies dans la documentation officielle d'Apache Kafka pour les configurations du producteur et les configurations du consommateur.

Realtime Compute for Apache Flink supprime le préfixe properties. et transmet les paires clé-valeur restantes au client Kafka sous-jacent. Par exemple, vous pouvez définir 'properties.allow.auto.create.topics' = 'false' pour désactiver la création automatique de topics.

Le connecteur Kafka écrase ces options ; vous ne pouvez donc pas les configurer de cette manière :

  • key.deserializer

  • value.deserializer

format

Format de sérialisation et de désérialisation de la valeur d'un message Kafka.

String

Non

Formats pris en charge :

  • csv

  • json

  • avro

  • debezium-json

  • canal-json

  • maxwell-json

  • avro-confluent

  • raw

Remarque

Pour plus d'informations, consultez les options de format.

key.format

Format de sérialisation et de désérialisation de la clé d'un message Kafka.

String

Non

Formats pris en charge :

  • csv

  • json

  • avro

  • debezium-json

  • canal-json

  • maxwell-json

  • avro-confluent

  • raw

Remarque

Lorsque vous utilisez cette configuration, key.options est obligatoire.

key.fields

Champs du schéma de table à utiliser comme clé du message Kafka.

String

Non

Séparez plusieurs noms de champs par des points-virgules (;). Par exemple, 'field1;field2'.

key.fields-prefix

Préfixe personnalisé pour tous les champs de clé afin d'éviter les conflits de nom avec les champs de valeur.

String

Non

Ce préfixe permet de distinguer les champs de clé des champs de valeur. Il est supprimé avant la sérialisation de la clé ou après sa désérialisation.

Remarque

Si vous utilisez cette option, value.fields-include doit être défini sur EXCEPT_KEY.

value.format

Format de sérialisation et de désérialisation de la valeur d'un message Kafka.

String

Non

Cette configuration équivaut à format. Vous ne pouvez définir que l'une des options format ou value.format. Si les deux sont configurées, value.format remplace format.

value.fields-include

Définit si les champs de clé sont inclus dans le format de valeur.

String

Non

ALL

Valeurs valides :

  • ALL : La valeur du message Kafka inclut toutes les colonnes de la table.

  • EXCEPT_KEY : La valeur du message Kafka inclut toutes les colonnes de la table, à l'exception de celles définies dans key.fields.

  • Table source

    Option

    Description

    Type

    Obligatoire

    Valeur par défaut

    Remarques

    topic

    Topic ou topics à lire.

    String

    Non

    Pour s'abonner à plusieurs topics, séparez leurs noms par des points-virgules (;), par exemple 'topic-1;topic-2'.

    Remarque

    Vous pouvez spécifier cette option ou topic-pattern, mais pas les deux simultanément.

    topic-pattern

    Expression régulière correspondant aux topics auxquels s'abonner. Le consommateur s'abonne à tous les topics dont les noms correspondent à ce modèle.

    String

    Non

    Exemples :

    • user_event_.* : Correspond à tous les topics préfixés par user_event_.

    • prod\.logs\..* : Correspond aux topics préfixés par prod.logs. (le caractère . doit être échappé).

    Remarque

    Vous pouvez spécifier cette option ou topic, mais pas les deux simultanément.

    properties.group.id

    ID du groupe de consommateurs source Kafka.

    String

    Non

    KafkaSource-{Nom-de-la-table-source}

    Si vous utilisez un ID de groupe de consommateurs pour la première fois, vous devez également définir properties.auto.offset.reset sur earliest ou latest afin de définir le décalage de démarrage initial.

    scan.startup.mode

    Décalage de démarrage du consommateur Kafka.

    String

    Non

    group-offsets

    Valeurs valides :

    • earliest-offset : Démarre la lecture à partir du décalage disponible le plus ancien.

    • latest-offset : Démarre la lecture à partir du dernier décalage.

    • group-offsets : Démarre la lecture à partir des décalages validés du properties.group.id spécifié.

    • timestamp : Démarre la lecture à partir du scan.startup.timestamp-millis spécifié.

    • specific-offsets : Démarre la lecture à partir des décalages spécifiés dans scan.startup.specific-offsets.

    Remarque

    Cette option s'applique uniquement lorsqu'un travail démarre sans état. Si un travail reprend à partir d'un point de contrôle, il lit à partir des décalages stockés dans l'état du point de contrôle.

    scan.startup.specific-offsets

    Décalage de démarrage par partition lorsque scan.startup.mode est défini sur specific-offsets.

    String

    Non

    Par exemple, partition:0,offset:42;partition:1,offset:300

    scan.startup.timestamp-millis

    Horodatage de démarrage en millisecondes lorsque scan.startup.mode est défini sur timestamp.

    Long

    Non

    L'unité est la milliseconde.

    scan.topic-partition-discovery.interval

    Intervalle de découverte des partitions.

    Duration

    Non

    5 minutes

    Le connecteur découvre et lit périodiquement les nouvelles partitions. Lorsque vous utilisez topic-pattern, le connecteur découvre également les nouveaux topics correspondant au modèle. Définissez l'intervalle sur une valeur non positive pour désactiver cette fonctionnalité.

    Remarque

    Dans Ververica Runtime (VVR) 6,0.x, la découverte dynamique des partitions est désactivée par défaut. À partir de VVR 8.0, cette fonctionnalité est activée par défaut avec un intervalle de découverte de 5 minutes.

    scan.header-filter

    Filtre les messages en fonction des en-têtes des messages Kafka.

    String

    Non

    Une clé d'en-tête et sa valeur sont séparées par deux-points (:). Plusieurs conditions d'en-tête sont reliées à l'aide d'opérateurs logiques (& et |). L'opérateur logique NON (!) est également pris en charge. Par exemple, depart:toy|depart:book&!env:test conserve les données Kafka si l'en-tête contient depart=toy ou depart=book et ne contient pas env=test.

    Remarque
    • Cette option est prise en charge uniquement dans Ververica Runtime (VVR) 8.0.6 et versions ultérieures.

    • Les parenthèses dans les expressions ne sont pas prises en charge.

    • Les opérations logiques sont évaluées de gauche à droite.

    • Les valeurs d'en-tête sont converties en chaînes UTF-8 pour la comparaison.

    scan.check.duplicated.group.id

    Vérifie si un autre consommateur actif utilise déjà le properties.group.id.

    Boolean

    Non

    false

    Valeurs valides :

    • true : Avant de démarrer le travail, le système vérifie l'existence d'un groupe de consommateurs en double. S'il en trouve un, le travail échoue pour éviter les conflits.

    • false : Démarre le travail sans vérifier les conflits.

    Remarque

    Cette option est prise en charge uniquement dans Ververica Runtime (VVR) 6.0.4 et versions ultérieures.

  • Table sink

    Option

    Description

    Type

    Obligatoire

    Valeur par défaut

    Remarques

    topic

    Topic cible.

    String

    Oui

    sink.partitioner

    Mappe les enregistrements des instances de sink parallèles aux partitions Kafka.

    String

    Non

    default

    Valeurs valides :

    • default : Utilise le répartiteur de partitions Kafka par défaut.

    • fixed : Chaque instance de sink parallèle écrit dans une partition Kafka fixe.

    • round-robin : Les enregistrements sont distribués aux partitions selon un algorithme tournant (round-robin).

    • Répartiteur personnalisé : Pour utiliser un répartiteur personnalisé, fournissez le nom de classe complet d'une sous-classe FlinkKafkaPartitioner, par exemple org.mycompany.MyPartitioner.

    sink.delivery-guarantee

    Garantie de livraison du sink.

    String

    Non

    at-least-once

    Valeurs valides :

    • none : Ne fournit aucune garantie. Les enregistrements peuvent être perdus ou dupliqués.

    • at-least-once : Garantit qu'aucun enregistrement n'est perdu, mais ils peuvent être dupliqués.

    • exactly-once : Utilise les transactions Kafka pour fournir une sémantique exactly-once, garantissant que les enregistrements ne sont ni perdus ni dupliqués.

    Remarque

    Lors de l'utilisation de la sémantique exactly-once, vous devez également spécifier sink.transactional-id-prefix.

    sink.transactional-id-prefix

    Préfixe d'ID de transaction. Obligatoire lorsque sink.delivery-guarantee est défini sur exactly-once.

    String

    Oui, si sink.delivery-guarantee est défini sur exactly-once

    Requis uniquement lorsque sink.delivery-guarantee est défini sur exactly-once.

    sink.parallelism

    Parallélisme de l'opérateur de sink.

    Integer

    Non

    Par défaut, le framework détermine le parallélisme en fonction des opérateurs en amont.

  • Utilisez le connecteur Kafka comme source, sink ou destination Flink CDC dans Realtime Compute for Apache Flink.

    Vue d'ensemble

    Apache Kafka est une plateforme open source distribuée de streaming d'événements, largement utilisée pour le traitement de données haute performance, l'analytique en continu et l'intégration de données. Le connecteur Kafka pour Realtime Compute for Apache Flink utilise le client Apache Kafka open source pour offrir un débit de données élevé, prendre en charge la lecture et l'écriture de plusieurs formats de données, et proposer une sémantique exactly-once.

    Catégorie

    Description

    Types pris en charge

    Source SQL, sink

    Source Flink CDC, sink

    Source DataStream, sink

    Mode d'exécution

    Streaming

    Formats de données

    Formats de données pris en charge

    • CSV

    • JSON

    • Apache Avro

    • Confluent Avro

    • Debezium JSON

    • Canal JSON

    • Maxwell JSON

    • Raw

    • Protobuf

    Remarque
    • Le format de données Protobuf intégré est pris en charge uniquement pour Ververica Runtime (VVR) 8.0.9 et versions ultérieures.

    • Chaque format de données pris en charge dispose de paramètres correspondants pouvant être spécifiés dans la clause WITH. Pour plus d'informations, consultez les formats.

    Métriques

    Métriques

    • Table source

      • numRecordsIn

      • numRecordsInPerSecond

      • numBytesIn

      • numBytesInPerSecond

      • currentEmitEventTimeLag

      • currentFetchEventTimeLag

      • sourceIdleTime

      • pendingRecords

    • Table sink

      • numRecordsOut

      • numRecordsOutPerSecond

      • numBytesOut

      • numBytesOutPerSecond

      • currentSendTime

    Remarque

    Pour plus d'informations sur les métriques, consultez les métriques.

    Types d'API

    SQL, DataStream, Flink CDC

    Mise à jour/suppression du sink

    Le connecteur prend uniquement en charge l'ajout de données à une table sink. Les mises à jour et les suppressions ne sont pas prises en charge.

    Remarque

    Pour plus d'informations sur la mise à jour ou la suppression de données dans une table sink, consultez Upsert Kafka.

    Prérequis

    Avant de commencer, vérifiez que vous remplissez les prérequis pour votre type de cluster Kafka :

    • Connexion à un cluster ApsaraMQ for Kafka

      • La version du cluster Kafka est 0,11 ou ultérieure.

      • Vous avez créé un cluster ApsaraMQ for Kafka. Pour plus d'informations, consultez l'étape Étape 3 : Créer des ressources.

      • L'espace de travail Flink et le cluster Kafka se trouvent dans le même Virtual Private Cloud (VPC), et vous avez ajouté le bloc CIDR de l'espace de travail Flink à la liste d'autorisation ApsaraMQ for Kafka. Pour plus d'informations, consultez la section Configurer les listes d'autorisation.

      Important

      Limitations liées à l'écriture de données dans ApsaraMQ for Kafka :

      • ApsaraMQ for Kafka ne prend pas en charge le format de compression Zstandard (zstd) pour les écritures.

      • ApsaraMQ for Kafka ne prend pas en charge les écritures idempotentes ou transactionnelles, ce qui empêche l'utilisation de la sémantique exactly-once fournie par les tables sink Kafka. À partir de Ververica Runtime (VVR) 8.0.0, le connecteur Kafka utilise le client Kafka 3.x, où la propriété properties.enable.idempotence est définie par défaut sur true. Par conséquent, pour éviter les échecs d'écriture lors de l'utilisation de Ververica Runtime (VVR) 8.0.0 ou version ultérieure pour écrire dans ApsaraMQ for Kafka, vous devez ajouter la configuration properties.enable.idempotence=false à la définition de votre table sink. Pour une comparaison des moteurs de stockage et des limitations de fonctionnalités pour ApsaraMQ for Kafka, consultez la section Comparaison entre les moteurs de stockage.

    • Connexion à un cluster Apache Kafka autogéré

      • La version du cluster Apache Kafka autogéré est 0,11 ou ultérieure.

      • L'espace de travail Flink dispose d'une connectivité réseau avec le cluster Apache Kafka autogéré. Pour plus de détails sur la connexion à un cluster via Internet public, consultez la section FAQ sur la connectivité réseau.

      • Seules les options de configuration client pour Apache Kafka version 2.8 sont prises en charge. Pour plus d'informations, consultez la documentation Apache Kafka relative aux configurations du consommateur et aux configurations du producteur.

    Remarques

    Les écritures transactionnelles ne sont pas recommandées en raison de limitations de conception connues dans Apache Flink et Apache Kafka. Lorsque vous définissez sink.delivery-guarantee = 'exactly-once', le connecteur Kafka active les écritures transactionnelles, avec les problèmes connus suivants :

    • Chaque point de contrôle génère un nouvel ID de transaction. Si l'intervalle de point de contrôle est trop court, le grand nombre d'ID de transaction peut entraîner une saturation de la mémoire du coordinateur du cluster Kafka, compromettant ainsi la stabilité du cluster.

    • Chaque transaction crée une nouvelle instance de producteur. Si trop de transactions sont validées simultanément, le TaskManager peut saturer sa mémoire, déstabilisant ainsi le travail Apache Flink.

    • Si plusieurs travaux Apache Flink utilisent le même sink.transactional-id-prefix, les ID de transaction générés peuvent entrer en conflit. Lorsqu'une opération d'écriture échoue dans un travail, cela peut empêcher l'avancement du Log Start Offset (LSO) d'une partition Apache Kafka. Cela affecte tous les consommateurs de cette partition.

    Si vous avez besoin d'une sémantique exactly-once, utilisez le connecteur Upsert Kafka pour écrire dans une table à clé primaire, garantissant ainsi l'idempotence. Si vous devez utiliser des écritures transactionnelles, consultez les remarques d'utilisation de la sémantique exactly-once.

    Dépannage de la connectivité réseau

    Une erreur Timed out waiting for a node assignment lors du démarrage d'un travail Realtime Compute for Apache Flink indique généralement un problème de connectivité réseau entre Realtime Compute for Apache Flink et le cluster Kafka.

    Un client Kafka se connecte aux courtiers de la manière suivante :

    1. Le client utilise les adresses spécifiées dans bootstrap.servers pour établir une connexion initiale au cluster Kafka.

    2. Le cluster Kafka renvoie les métadonnées de chaque courtier, y compris leurs endpoints.

    3. Le client utilise ensuite ces endpoints pour se connecter aux courtiers afin de lire ou d'écrire des données.

    Même si les adresses bootstrap.servers sont accessibles, le client ne peut pas lire ou écrire de données si Kafka renvoie des endpoints de courtier incorrects. Ce problème survient souvent dans les architectures réseau utilisant un proxy, une redirection de port ou une ligne louée.

    Étapes de dépannage

    ApsaraMQ for Kafka

    1. Confirmez le type d'endpoint

      • Endpoint par défaut (réseau interne)

      • Endpoint SASL (réseau interne avec authentification)

      • Endpoint public (nécessite une demande distincte)

      Utilisez la fonctionnalité Network Probe dans la console de développement Realtime Compute for Apache Flink pour écarter les problèmes de connectivité avec l'adresse bootstrap.servers.

    2. Vérifiez les groupes de sécurité et les listes d'autorisation

      Ajoutez le bloc CIDR de l'espace de travail Realtime Compute for Apache Flink à la liste d'autorisation de votre instance Kafka. Pour plus d'informations, consultez les sections Afficher le bloc CIDR du VPC et Configurer une liste d'autorisation.

    3. Vérifiez la configuration SASL (si activée)

      Si vous utilisez un endpoint SASL_SSL, assurez-vous que les mécanismes JAAS, SSL et SASL sont correctement configurés dans votre travail Realtime Compute for Apache Flink. Sans une authentification appropriée, la connexion peut échouer pendant la phase de handshake, ce qui peut également se manifester par un délai d'expiration. Pour plus d'informations, consultez la section Sécurité et authentification.

    Self-managed Kafka

    1. Utilisez la fonctionnalité Network Probe

      Cette fonctionnalité vous aide à écarter les problèmes de connectivité avec l'adresse bootstrap.servers et à vérifier que l'endpoint interne ou public correct est utilisé.

    2. Vérifiez les groupes de sécurité et les listes d'autorisation

      • Le groupe de sécurité de l'instance Elastic Compute Service (ECS) doit autoriser le trafic entrant sur le port de l'endpoint Kafka, généralement 9092 ou 9093.

      • Assurez-vous que tout pare-feu sur l'instance ECS autorise le trafic provenant du VPC de votre espace de travail Realtime Compute for Apache Flink. Pour plus d'informations, consultez la section Afficher le bloc CIDR du VPC.

    3. Vérifiez la configuration

      1. Utilisez l'outil zkCli.sh ou zookeeper-shell.sh pour vous connecter au cluster ZooKeeper utilisé par Kafka.

      2. Exécutez une commande pour obtenir les métadonnées du courtier. Par exemple, exécutez get /brokers/ids/0. Dans le champ endpoints de la réponse, recherchez l'adresse que Kafka annonce aux clients.

        # bin/zookeeper-shell.sh localhost:2181
        Connecting to localhost:2181
        Welcome to ZooKeeper!
        JLine support is disabled
        
        WATCHER::
        
        WatchedEvent state:SyncConnected type:None path:null
        get /brokers/ids/0
        {"listener_security_protocol_map":{"PLAINTEXT":"PLAINTEXT"},"endpoints":["PLAINTEXT://116.62.xxx:9092"],"jmx_port":-1,"host":"116.62.xxx","timestamp":"1614840078030","port":9092,"version":4}
        
      3. Utilisez la fonctionnalité Network Probe dans la console de développement Realtime Compute for Apache Flink pour tester l'accessibilité de cette adresse.

        Remarque
        • Si l'adresse n'est pas accessible, contactez vos administrateurs Kafka pour vérifier et corriger les configurations listeners et advertised.listeners afin de garantir que l'adresse annoncée est accessible depuis Realtime Compute for Apache Flink.

        • Pour plus d'informations sur les connexions des clients Kafka, consultez la section Dépannage de la connectivité.

    4. Vérifiez la configuration SASL (si activée)

      Si vous utilisez un endpoint SASL_SSL, assurez-vous que les mécanismes JAAS, SSL et SASL sont correctement configurés dans votre travail Realtime Compute for Apache Flink. Sans une authentification appropriée, la connexion peut échouer pendant la phase de handshake, ce qui peut également se manifester par un délai d'expiration. Pour plus d'informations, consultez la section Sécurité et authentification.

    SQL

    Utilisez le connecteur Kafka comme table source ou table sink dans les travaux SQL.

    Syntaxe

    CREATE TABLE KafkaTable (
      `user_id` BIGINT,
      `item_id` BIGINT,
      `behavior` STRING,
      `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'user_behavior',
      'properties.bootstrap.servers' = 'localhost:9092',
      'properties.group.id' = 'testGroup',
      'scan.startup.mode' = 'earliest-offset',
      'format' = 'csv'
    )

    Colonnes de métadonnées

    Définissez des colonnes de métadonnées dans une table source ou sink pour accéder aux métadonnées des messages Kafka. Par exemple, lorsque vous vous abonnez à plusieurs topics, une colonne de métadonnées peut identifier le topic d'origine de chaque enregistrement.

    CREATE TABLE kafka_source (
      -- Read the message topic as the `record_topic` column
      `record_topic` STRING NOT NULL METADATA FROM 'topic' VIRTUAL,
      -- Read the timestamp from the ConsumerRecord as the `ts` column
      `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL,
      -- Read the message offset as the `record_offset` column
      `record_offset` BIGINT NOT NULL METADATA FROM 'offset' VIRTUAL,
      ...
    ) WITH (
      'connector' = 'kafka',
      ...
    );
    
    CREATE TABLE kafka_sink (
      -- Write the timestamp from the `ts` column as the ProducerRecord's timestamp to Kafka
      `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL,
      ...
    ) WITH (
      'connector' = 'kafka',
      ...
    );

    Le tableau suivant répertorie les colonnes de métadonnées prises en charge par les tables source et sink Kafka.

    Clé

    Type

    Description

    Portée

    topic

    STRING NOT NULL METADATA VIRTUAL

    Topic du message.

    Table source

    partition

    INT NOT NULL METADATA VIRTUAL

    ID de partition du message.

    Table source

    headers

    MAP<STRING, BYTES> NOT NULL METADATA VIRTUAL

    En-têtes du message.

    Table source et table sink

    leader-epoch

    INT NOT NULL METADATA VIRTUAL

    Leader-epoch du message.

    Table source

    offset

    BIGINT NOT NULL METADATA VIRTUAL

    Décalage du message.

    Table source

    timestamp

    TIMESTAMP(3) WITH LOCAL TIME ZONE NOT NULL METADATA VIRTUAL

    Horodatage du message.

    Table source et table sink

    timestamp-type

    STRING NOT NULL METADATA VIRTUAL

    Type d'horodatage du message. Les valeurs valides sont :

    • NoTimestampType : Aucun horodatage n'est défini dans le message.

    • CreateTime : Heure de création du message.

    • LogAppendTime : Heure d'ajout du message au journal du courtier Kafka.

    Table source

    __raw_key__

    STRING NOT NULL METADATA VIRTUAL

    Clé brute du message.

    Table source et table sink

    Remarque

    Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.4 et versions ultérieures.

    __raw_value__

    STRING NOT NULL METADATA VIRTUAL

    Valeur brute du message.

    Table source et table sink

    Remarque

    Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.4 et versions ultérieures.

    Sécurité et authentification

    Si le cluster Kafka nécessite une connexion sécurisée ou une authentification, préfixez les configurations de sécurité et d'authentification pertinentes par properties. et définissez-les dans le paramètre WITH. L'exemple suivant configure une table Kafka pour utiliser PLAIN comme mécanisme SASL avec une configuration JAAS.

    CREATE TABLE KafkaTable (
      `user_id` BIGINT,
      `item_id` BIGINT,
      `behavior` STRING,
      `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
    ) WITH (
      'connector' = 'kafka',
      ...
      'properties.security.protocol' = 'SASL_PLAINTEXT',
      'properties.sasl.mechanism' = 'PLAIN',
      'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="username" password="password";'
    )

    L'exemple suivant montre comment utiliser SASL_SSL comme protocole de sécurité et SCRAM-SHA-256 comme mécanisme SASL.

    CREATE TABLE KafkaTable (
      `user_id` BIGINT,
      `item_id` BIGINT,
      `behavior` STRING,
      `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
    ) WITH (
      'connector' = 'kafka',
      ...
      'properties.security.protocol' = 'SASL_SSL',
      /* SSL configuration */
      /* Path to the truststore for the server's CA certificate. */
      /* Files uploaded using Artifacts are stored in the /flink/usrlib/ directory. */
      'properties.ssl.truststore.location' = '/flink/usrlib/kafka.client.truststore.jks',
      'properties.ssl.truststore.password' = 'test1234',
      /* If client authentication is required, you must also configure the path to the keystore (private key). */
      'properties.ssl.keystore.location' = '/flink/usrlib/kafka.client.keystore.jks',
      'properties.ssl.keystore.password' = 'test1234',
      /* The algorithm used to verify the server hostname. An empty string disables hostname verification. */
      'properties.ssl.endpoint.identification.algorithm' = '',
      /* SASL configuration */
      /* Set the SASL mechanism to SCRAM-SHA-256. */
      'properties.sasl.mechanism' = 'SCRAM-SHA-256',
      /* Configure JAAS. */
      'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModule required username="username" password="password";'
    )

    Vous pouvez utiliser la fonctionnalité Artifacts de la console Realtime Compute for Apache Flink pour télécharger le certificat CA et la clé privée mentionnés dans l'exemple. Les fichiers téléchargés sont stockés dans le répertoire /flink/usrlib. Pour utiliser un fichier de certificat CA nommé my-truststore.jks, vous pouvez définir la propriété 'properties.ssl.truststore.location' dans la clause WITH de l'une des deux manières suivantes :

    • Définissez 'properties.ssl.truststore.location' = '/flink/usrlib/my-truststore.jks'. Cette méthode évite le téléchargement dynamique de fichiers depuis Object Storage Service (OSS) au moment de l'exécution, mais elle ne prend pas en charge le mode Debug.

    • Si la version du moteur Realtime Compute est VVR 11.5 ou ultérieure, vous pouvez configurer properties.ssl.truststore.location et properties.ssl.keystore.location avec un chemin OSS absolu. Le format du chemin de fichier est oss://flink-fullymanaged-<ID de l'espace de travail>/artifacts/namespaces/<Nom de l'espace de noms>/<nom du fichier>. Cette méthode télécharge dynamiquement les fichiers OSS lors de l'exécution de Flink et prend en charge le mode Debug.

    Remarque
    • Vérifiez votre configuration : Les exemples de cette rubrique présentent des configurations courantes. Avant de configurer le connecteur Kafka, contactez votre équipe d'exploitation et de maintenance Kafka pour obtenir les paramètres de sécurité et d'authentification corrects.

    • Échappement : Contrairement à Apache Flink natif, l'éditeur SQL de Realtime Compute for Apache Flink échappe par défaut les guillemets doubles ("). Par conséquent, vous n'avez pas besoin d'ajouter des barres obliques inverses (\) pour échapper les guillemets doubles utilisés pour le nom d'utilisateur et le mot de passe dans l'option properties.sasl.jaas.config.

    Décalage de démarrage de la table source

    Mode de démarrage

    Vous pouvez configurer l'option scan.startup.mode pour spécifier le décalage à partir duquel une table source Kafka commence à lire les données. Les valeurs valides incluent :

    • earliest-offset : Démarre la lecture à partir du décalage le plus ancien.

    • latest-offset : Démarre la lecture à partir du dernier décalage.

    • group-offsets : Démarre la lecture à partir des décalages validés pour le groupe de consommateurs spécifié dans properties.group.id.

    • timestamp : Démarre la lecture à partir du premier message dont l'horodatage est supérieur ou égal à la valeur spécifiée dans scan.startup.timestamp-millis.

    • specific-offsets : Démarre la lecture à partir des décalages de partition spécifiques spécifiés dans scan.startup.specific-offsets.

    Remarque
    • Si vous ne spécifiez pas de mode de démarrage, la valeur par défaut est 'group-offsets'.

    • L'option scan.startup.mode s'applique uniquement aux travaux sans état. Lorsqu'un travail avec état démarre, il consomme toujours à partir des décalages stockés dans son état.

    Exemple :

    CREATE TEMPORARY TABLE kafka_source (
      ...
    ) WITH (
      'connector' = 'kafka',
      ...
      -- Consume from the earliest offset.
      'scan.startup.mode' = 'earliest-offset',
      -- Consume from the latest offset.
      'scan.startup.mode' = 'latest-offset',
      -- Consume from the committed offsets of the consumer group "my-group".
      'properties.group.id' = 'my-group',
      'scan.startup.mode' = 'group-offsets',
      'properties.auto.offset.reset' = 'earliest', -- If "my-group" is used for the first time, consumption starts from the earliest offset.
      'properties.auto.offset.reset' = 'latest', -- If "my-group" is used for the first time, consumption starts from the latest offset.
      -- Consume from the specified timestamp in milliseconds: 1655395200000.
      'scan.startup.mode' = 'timestamp',
      'scan.startup.timestamp-millis' = '1655395200000',
      -- Consume from specific offsets.
      'scan.startup.mode' = 'specific-offsets',
      'scan.startup.specific-offsets' = 'partition:0,offset:42;partition:1,offset:300'
    );

    Priorité du décalage de démarrage

    Le décalage de démarrage de la table source est déterminé par les règles suivantes, par ordre de priorité :

    Priorité (de la plus élevée à la plus basse)

    Le décalage stocké dans un point de contrôle ou un savepoint.

    L'heure de démarrage sélectionnée dans la console Realtime Compute for Apache Flink lors du démarrage du travail.

    Le décalage de démarrage spécifié par scan.startup.mode dans la clause WITH.

    Si scan.startup.mode n'est pas spécifié, group-offsets est utilisé pour démarrer la consommation à partir des décalages du groupe de consommateurs correspondant.

    Si le décalage déterminé par l'une de ces étapes est invalide, par exemple parce qu'il a expiré ou qu'un problème est survenu dans le cluster Kafka, le système réinitialise le décalage selon la politique spécifiée dans properties.auto.offset.reset. Si cette option n'est pas configurée, le système lève une exception nécessitant une intervention utilisateur.

    Un scénario courant implique le démarrage de la consommation avec un nouvel ID de groupe de consommateurs. La table source interroge d'abord le cluster Kafka pour obtenir les décalages validés de ce groupe. Comme l'ID de groupe est nouveau, aucun décalage valide n'est trouvé. Par conséquent, le système réinitialise le décalage selon la politique spécifiée dans properties.auto.offset.reset. Ainsi, lors de la consommation avec un nouvel ID de groupe, vous devez configurer l'option properties.auto.offset.reset.

    Validation des décalages source

    La table source Kafka valide son décalage de consommateur auprès du cluster Kafka uniquement après un point de contrôle réussi ; ainsi, un long intervalle de point de contrôle entraîne un retard du décalage validé. La table source stocke la progression réelle de la lecture dans l'état du point de contrôle, que le système utilise pour la récupération après incident. Les décalages validés servent uniquement de moniteur de progression et ne sont pas utilisés pour la récupération ; par conséquent, les échecs de validation n'affectent pas l'exactitude des données.

    Répartiteur de partitions de sink personnalisé

    Si la stratégie de partitionnement intégrée de Kafka ne répond pas à vos besoins, vous pouvez implémenter un répartiteur personnalisé en étendant la classe FlinkKafkaPartitioner. Une fois le développement terminé, compilez votre code dans un package JAR et téléchargez-le à l'aide de la fonctionnalité Artifacts dans la console Realtime Compute. Après le téléchargement et la référence du package JAR, définissez le paramètre sink.partitioner dans la clause WITH avec le nom de classe complet de votre répartiteur, par exemple org.mycompany.MyPartitioner.

    Kafka, Upsert Kafka et catalogue Kafka JSON

    Kafka est une plateforme de streaming d'événements en ajout seul qui ne prend pas en charge les mises à jour ou les suppressions de données. En SQL streaming, une table sink Kafka standard ne peut pas gérer les données Change Data Capture (CDC) en amont ni la logique de rétractation des opérateurs tels que l'agrégation et la jointure. Si vous devez écrire des données contenant des modifications ou des rétractations, utilisez une table sink Upsert Kafka.

    Pour simplifier la synchronisation par lots des données Change Data Capture (CDC) d'une ou plusieurs tables de base de données en amont vers Kafka, vous pouvez utiliser un catalogue Kafka JSON. Si les données stockées dans Kafka sont au format JSON, un catalogue Kafka JSON vous permet d'ignorer l'étape de définition du schéma et des paramètres WITH. Pour plus de détails, consultez la section Gérer les catalogues Kafka JSON.

    Exemples

    Exemple 1 : Lire et écrire dans Kafka

    Cet exemple lit les données d'un topic source Kafka et les écrit dans un topic sink. Les données sont au format CSV.

    CREATE TEMPORARY TABLE kafka_source (
      id INT,
      name STRING,
      age INT
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'source',
      'properties.bootstrap.servers' = '<yourKafkaBrokers>',
      'properties.group.id' = '<yourKafkaConsumerGroupId>',
      'format' = 'csv'
    );
    
    CREATE TEMPORARY TABLE kafka_sink (
      id INT,
      name STRING,
      age INT
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'sink',
      'properties.bootstrap.servers' = '<yourKafkaBrokers>',
      'properties.group.id' = '<yourKafkaConsumerGroupId>',
      'format' = 'csv'
    );
    
    INSERT INTO kafka_sink SELECT id, name, age FROM kafka_source;

    Exemple 2 : Synchroniser le schéma de table et les données

    Vous pouvez utiliser le connecteur Kafka pour synchroniser en temps réel les messages d'un topic Kafka vers Hologres. Pour éviter les messages en double dans Hologres lors d'un basculement, vous pouvez utiliser le décalage et l'ID de partition des messages Kafka comme clé primaire composite.

    CREATE TEMPORARY TABLE kafkaTable (
      `offset` INT NOT NULL METADATA,
      `part` BIGINT NOT NULL METADATA FROM 'partition',
      PRIMARY KEY (`part`, `offset`) NOT ENFORCED
    ) WITH (
      'connector' = 'kafka',
      'properties.bootstrap.servers' = '<yourKafkaBrokers>',
      'topic' = 'kafka_evolution_demo',
      'scan.startup.mode' = 'earliest-offset',
      'format' = 'json',
      'json.infer-schema.flatten-nested-columns.enable' = 'true'
        -- Optional. Flattens all nested columns.
    );
    
    CREATE TABLE IF NOT EXISTS hologres.kafka.`sync_kafka`
    WITH (
      'connector' = 'hologres'
    ) AS TABLE vvp.`default`.kafkaTable;

    Exemple 3 : Synchroniser les clés et les valeurs Kafka

    Si la clé d'un message Kafka contient des informations pertinentes, vous pouvez synchroniser à la fois la clé et la valeur.

    CREATE TEMPORARY TABLE kafkaTable (
      `key_id` INT NOT NULL,
      `val_name` VARCHAR(200)
    ) WITH (
      'connector' = 'kafka',
      'properties.bootstrap.servers' = '<yourKafkaBrokers>',
      'topic' = 'kafka_evolution_demo',
      'scan.startup.mode' = 'earliest-offset',
      'key.format' = 'json',
      'value.format' = 'json',
      'key.fields' = 'key_id',
      'key.fields-prefix' = 'key_',
      'value.fields-prefix' = 'val_',
      'value.fields-include' = 'EXCEPT_KEY'
    );
    
    CREATE TABLE IF NOT EXISTS hologres.kafka.`sync_kafka`(
    WITH (
      'connector' = 'hologres'
    ) AS TABLE vvp.`default`.kafkaTable;
    Remarque

    Les clés de message Kafka ne prennent pas en charge l'évolution du schéma ni l'analyse automatique des types. Vous devez déclarer le schéma manuellement.

    Exemple 4 : Synchroniser les données et effectuer un calcul

    Lors de la synchronisation des données de Kafka vers Hologres, vous pouvez avoir besoin de transformations légères.

    CREATE TEMPORARY TABLE kafkaTable (
      `distinct_id` INT NOT NULL,
      `properties` STRING,
      `timestamp` TIMESTAMP_LTZ METADATA,
      `date` AS CAST(`timestamp` AS DATE)
    ) WITH (
      'connector' = 'kafka',
      'properties.bootstrap.servers' = '<yourKafkaBrokers>',
      'topic' = 'kafka_evolution_demo',
      'scan.startup.mode' = 'earliest-offset',
      'key.format' = 'json',
      'value.format' = 'json',
      'key.fields' = 'key_id',
      'key.fields-prefix' = 'key_'
    );
    
    CREATE TABLE IF NOT EXISTS hologres.kafka.`sync_kafka` WITH (
       'connector' = 'hologres'
    ) AS TABLE vvp.`default`.kafkaTable
    ADD COLUMN
      `order_id` AS COALESCE(JSON_VALUE(`properties`, '$.order_id'), 'default');
    --Use COALESCE to handle null values.

    Exemple 5 : Analyser un JSON imbriqué

    Voici un exemple de message JSON :

    {
      "id": 101,
      "name": "VVP",
      "properties": {
        "owner": "Alibaba Cloud",
        "engine": "Flink"
      }
    }

    Pour éviter d'utiliser des fonctions telles que JSON_VALUE(payload, '$.properties.owner') pour analyser les champs, vous pouvez définir directement la structure dans le DDL Source :

    CREATE TEMPORARY TABLE kafka_source (
      id          VARCHAR,
      `name`      VARCHAR,
      properties  ROW<`owner` STRING, engine STRING>
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'xxx',
      'properties.bootstrap.servers' = 'xxx',
      'scan.startup.mode' = 'earliest-offset',
      'format' = 'json'
    );

    Avec cette approche, Flink analyse le JSON en champs structurés lors de la phase de lecture. Les requêtes SQL ultérieures peuvent référencer directement properties.owner sans appels de fonction supplémentaires, ce qui améliore les performances globales.

    API DataStream

    Important

    Pour lire ou écrire des données avec l'API DataStream, utilisez le connecteur DataStream correspondant pour vous connecter à Realtime Compute for Apache Flink. Pour plus d'informations sur la configuration d'un connecteur DataStream, consultez la section Intégrer des connecteurs DataStream.

    • Créer une source Kafka

      La source Kafka fournit une classe de générateur pour créer une instance de source Kafka. Le code exemple suivant crée une source Kafka qui consomme des données à partir du décalage le plus ancien du topic input-topic. Le groupe de consommateurs est my-group et la valeur du message Kafka est désérialisée sous forme de chaîne.

      Java

      KafkaSource<String> source = KafkaSource.<String>builder()
          .setBootstrapServers(brokers)
          .setTopics("input-topic")
          .setGroupId("my-group")
          .setStartingOffsets(OffsetsInitializer.earliest())
          .setValueOnlyDeserializer(new SimpleStringSchema())
          .build();
      
      env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");

      Pour créer une source Kafka, vous devez spécifier les propriétés suivantes.

      Paramètre

      Description

      BootstrapServers

      Liste des adresses des courtiers Kafka. Définissez cette propriété en appelant la méthode setBootstrapServers(String).

      GroupId

      ID du groupe de consommateurs. Définissez cette propriété en appelant la méthode setGroupId(String).

      Topics ou partitions

      Topics ou partitions auxquels s'abonner. La source Kafka prend en charge les trois méthodes suivantes pour s'abonner à des topics ou des partitions :

      • S'abonne à toutes les partitions des topics d'une liste.

        KafkaSource.builder().setTopics("topic-a","topic-b")
      • Modèle de topic : s'abonne à toutes les partitions des topics dont les noms correspondent à l'expression régulière spécifiée.

        KafkaSource.builder().setTopicPattern("topic.*")
      • Liste des partitions, permettant de s'abonner à une partition spécifique.

        final HashSet<TopicPartition> partitionSet = new HashSet<>(Arrays.asList(
                new TopicPartition("topic-a", 0),    // Partition 0 du topic "topic-a"
                new TopicPartition("topic-b", 5)));  // Partition 5 du topic "topic-b"
        KafkaSource.builder().setPartitions(partitionSet)

      Désérialiseur

      Désérialiseur utilisé pour analyser les messages Kafka.

      Spécifiez le désérialiseur à l'aide de la méthode setDeserializer(KafkaRecordDeserializationSchema). KafkaRecordDeserializationSchema définit comment analyser un ConsumerRecord Kafka. Si vous devez uniquement analyser la valeur d'un message Kafka, vous pouvez utiliser l'une des méthodes suivantes :

      • Utilisez la méthode setValueOnlyDeserializer(DeserializationSchema) de la classe de générateur. DeserializationSchema définit comment analyser les données binaires de la valeur du message Kafka.

      • Utilisez une classe qui implémente l'interface Deserializer de Kafka. Par exemple, vous pouvez utiliser StringDeserializer pour analyser la valeur du message Kafka en une chaîne.

        import org.apache.kafka.common.serialization.StringDeserializer;
        
        KafkaSource.<String>builder()
                .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class));
      Remarque

      Pour analyser un ConsumerRecord complet, vous devez implémenter l'interface KafkaRecordDeserializationSchema.

      POM

      Le connecteur Kafka DataStream est disponible dans le dépôt central Maven.

      <dependency>
          <groupId>com.alibaba.ververica</groupId>
          <artifactId>ververica-connector-kafka</artifactId>
          <version>${vvr-version}</version>
      </dependency>

      Lors de l'utilisation du connecteur DataStream Kafka, tenez compte des propriétés suivantes :

      • Décalage de démarrage

        Une source Kafka spécifie son décalage de démarrage à l'aide d'un initialiseur de décalage (OffsetsInitializer). Les initialiseurs intégrés incluent :

        Initialiseur de décalage

        Code

        Démarre la consommation à partir du décalage le plus ancien.

        KafkaSource.builder().setStartingOffsets(OffsetsInitializer.earliest())

        Démarre la consommation à partir du dernier décalage.

        KafkaSource.builder().setStartingOffsets(OffsetsInitializer.latest())

        Démarre la consommation des données dont l'horodatage est supérieur ou égal à l'heure spécifiée. L'unité est la milliseconde.

        KafkaSource.builder().setStartingOffsets(OffsetsInitializer.timestamp(1592323200000L))

        Démarre la consommation à partir du décalage validé du groupe de consommateurs. Si aucun décalage validé n'existe, il utilise la stratégie de réinitialisation spécifiée (par exemple, le décalage le plus ancien).

        KafkaSource.builder().setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))

        La consommation démarre à partir du décalage validé par le groupe de consommateurs, et aucune politique de réinitialisation de décalage n'est spécifiée.

        KafkaSource.builder().setStartingOffsets(OffsetsInitializer.committedOffsets())

        Remarque
        • Si les initialiseurs intégrés ne répondent pas à vos besoins, vous pouvez implémenter un initialiseur de décalage personnalisé.

        • Si vous ne spécifiez pas d'initialiseur de décalage, la valeur par défaut est OffsetsInitializer.earliest().

      • Mode streaming et mode batch

        La source Kafka prend en charge à la fois le mode streaming et le mode batch. Par défaut, elle fonctionne en mode streaming, où le travail s'exécute indéfiniment jusqu'à ce qu'il échoue ou soit annulé. Pour configurer la source Kafka afin qu'elle s'exécute en mode batch, vous pouvez utiliser setBounded(OffsetsInitializer) pour spécifier un décalage d'arrêt. La source Kafka se termine lorsque toutes les partitions atteignent leurs décalages d'arrêt spécifiés.

        Remarque

        Une source Kafka en mode streaming n'a généralement pas de décalage d'arrêt. Toutefois, à des fins de test, vous pouvez utiliser setUnbounded(OffsetsInitializer) pour spécifier un décalage d'arrêt même en mode streaming. Notez les noms de méthode différents pour spécifier le décalage d'arrêt : setUnbounded pour le mode streaming et setBounded pour le mode batch.

      • Découverte dynamique des partitions

        Pour gérer le scaling des topics ou la création de nouveaux topics sans redémarrer le travail Flink, vous pouvez activer la découverte dynamique des partitions lors de l'abonnement à des topics par modèle. Cette fonctionnalité est désactivée par défaut et doit être explicitement activée :

        KafkaSource.builder()
            .setProperty("partition.discovery.interval.ms", "10000") // Discover new partitions every 10 seconds.
        Important

        La fonctionnalité de découverte dynamique des partitions dépend du mécanisme de mise à jour des métadonnées du cluster Kafka. Si le cluster Kafka ne met pas à jour les informations de partition en temps opportun, les nouvelles partitions pourraient ne pas être découvertes. Assurez-vous que la configuration partition.discovery.interval.ms du cluster Kafka correspond à votre scénario réel.

      • Heure d'événement et watermark

        Par défaut, la source Kafka utilise l'horodatage du message Kafka comme heure d'événement. Vous pouvez définir une stratégie de watermark personnalisée pour extraire l'heure d'événement du corps du message et émettre un watermark en aval.

        env.fromSource(kafkaSource, new CustomWatermarkStrategy(), "Kafka Source With Custom Watermark Strategy")

        Pour en savoir plus sur les stratégies de watermark personnalisées, consultez la section Génération de watermarks.

        Remarque

        Si une sous-tâche source est inactive (par exemple, lorsqu'une partition Kafka ne contient pas de nouvelles données ou que le parallélisme de la source est supérieur au nombre de partitions Kafka), le watermark de cette sous-tâche n'avancera pas. Cela peut bloquer les calculs de fenêtre en aval.

        Pour résoudre ce problème, envisagez les solutions suivantes :

        • Configurez un délai d'inactivité de la source : activez la propriété table.exec.source.idle-timeout pour marquer une source inactive comme temporairement inactive. Cela permet au watermark en aval d'avancer.

        • Définissez un parallélisme approprié : assurez-vous que le parallélisme de la source n'est pas supérieur au nombre de partitions Kafka.

      • Validation du décalage

        Paramètres

        • Général

          Paramètre

          Description

          Obligatoire

          Type

          Valeur par défaut

          Remarques

          type

          Type de source ou de puits.

          Oui

          String

          La valeur doit être kafka.

          name

          Nom de la source ou du puits.

          Non

          String

          Aucun

          properties.bootstrap.servers

          Adresses des brokers Kafka.

          Oui

          String

          Le format est host1:port1,host2:port2,host3:port3, séparé par des virgules (,).

          properties.*

          Propriétés de configuration pour le client Kafka.

          Non

          String

          Les clés de propriété doivent correspondre aux options valides définies dans la documentation officielle d'Apache Kafka pour les configurations du producteur et les configurations du consommateur.

          Realtime Compute for Apache Flink (VVR) supprime le préfixe properties. avant de transmettre les paires clé-valeur restantes au client Kafka sous-jacent. Par exemple, 'properties.allow.auto.create.topics' = 'false' permet de désactiver la création automatique de topics.

          key.format

          Format de sérialisation et de désérialisation de la clé du message Kafka.

          Non

          String

          • Pour la source, seul le format json est pris en charge.

          • Pour le puits, les valeurs valides sont :

            • csv

            • json

          Remarque

          Cette option n'est prise en charge qu'à partir de Realtime Compute for Apache Flink (VVR) 11.0.0.

          value.format

          Format de sérialisation et de désérialisation de la valeur du message Kafka.

          Non

          String

          debezium-json

          • Pour la source, les valeurs valides sont :

            • debezium-json

            • canal-json

            • json

          • Pour le puits, les valeurs valides sont :

            • debezium-json

            • canal-json

            • canal-protobuf

          Remarque
          • Les formats debezium-json et canal-json nécessitent Realtime Compute for Apache Flink (VVR) version 8.0.10 ou ultérieure.

          • Le format json nécessite Realtime Compute for Apache Flink (VVR) version 11.0.0 ou ultérieure.

        • Paramètres de la source

          Paramètre

          Description

          Obligatoire

          Type

          Valeur par défaut

          Remarques

          topic

          Topic ou topics à lire.

          Non

          String

          Pour s'abonner à plusieurs topics, séparez leurs noms par des points-virgules (;), par exemple topic-1;topic-2.

          Remarque

          Spécifiez ce paramètre ou topic-pattern, mais pas les deux.

          topic-pattern

          Expression régulière correspondant aux noms des topics auxquels s'abonner.

          Non

          String

          Exemples :

          • user_event_.* : Correspond à tous les topics préfixés par user_event_.

          • prod\.logs\..* : Correspond aux topics préfixés par prod.logs. (le caractère . doit être échappé).

          Remarque

          Spécifiez ce paramètre ou topic, mais pas les deux.

          properties.group.id

          ID du groupe de consommateurs.

          Non

          String

          Lorsque vous spécifiez un nouvel ID de groupe de consommateurs, vous devez définir le paramètre properties.auto.offset.reset sur earliest ou latest afin de définir l'offset de départ initial.

          scan.startup.mode

          Offset de démarrage du consommateur Kafka.

          Non

          String

          group-offsets

          Valeurs valides :

          • earliest-offset : Commence la lecture à partir du premier offset disponible.

          • latest-offset : Commence la lecture à partir du dernier offset.

          • group-offsets (valeur par défaut) : Commence la lecture à partir des offsets validés pour le properties.group.id spécifié.

          • timestamp : Commence la lecture à partir de l'horodatage spécifié par scan.startup.timestamp-millis.

          • specific-offsets : Commence la lecture à partir des offsets spécifiés par scan.startup.specific-offsets.

          Remarque

          Ce paramètre s'applique uniquement lors du démarrage d'un job sans état. Lorsqu'un job avec état démarre, il consomme toujours à partir des offsets stockés dans son état.

          scan.startup.specific-offsets

          Offset de démarrage par partition lorsque scan.startup.mode est défini sur specific-offsets.

          Non

          String

          Par exemple, partition:0,offset:42;partition:1,offset:300

          scan.startup.timestamp-millis

          Horodatage de démarrage en millisecondes lorsque scan.startup.mode est défini sur timestamp.

          Non

          Long

          L'unité est la milliseconde.

          scan.topic-partition-discovery.interval

          Intervalle de découverte dynamique des nouvelles partitions au sein des topics.

          Non

          Duration

          5 minutes

          Le connecteur découvre et lit périodiquement les nouvelles partitions. Lorsque vous utilisez topic-pattern, le connecteur découvre également les nouveaux topics correspondant au modèle. Pour désactiver la découverte, définissez cette valeur sur 0 ou moins.

          scan.check.duplicated.group.id

          Vérifie si le groupe de consommateurs spécifié par properties.group.id est dupliqué.

          Non

          Boolean

          false

          Valeurs valides :

          • true : Vérifie la présence d'un groupe de consommateurs dupliqué avant le démarrage du job. Si un doublon est trouvé, le job échoue.

          • false : Démarre le job sans vérifier les conflits.

          schema.inference.strategy

          Stratégie d'analyse du schéma.

          Non

          String

          continuous

          Valeurs valides :

          • continuous : Analyse le schéma de chaque enregistrement de données. Si les schémas sont incompatibles, le système déduit un schéma plus large et génère un événement de modification de schéma.

          • static : Effectue l'analyse du schéma une seule fois au démarrage du job. Les données sont ensuite analysées sur la base de ce schéma initial et aucun événement de modification de schéma n'est généré.

          Remarque

          scan.max.pre.fetch.records

          Nombre maximal de messages consommés par partition pour l'inférence initiale du schéma.

          Non

          Int

          50

          Avant le début du traitement des données, le système récupère et consomme le nombre spécifié de messages récents de chaque partition pour initialiser le schéma.

          key.fields-prefix

          Préfixe pour les noms de champs de la clé du message afin d'éviter les conflits de nom.

          Non

          String

          Par exemple, si ce paramètre est défini sur key_ et que la clé du message contient un champ nommé a, le nom du champ analysé devient key_a.

          Remarque

          La valeur de key.fields-prefix ne peut pas être un préfixe de la valeur de value.fields-prefix.

          value.fields-prefix

          Préfixe pour les noms de champs de la valeur du message afin d'éviter les conflits de nom.

          Non

          String

          Par exemple, si ce paramètre est défini sur value_ et que la valeur du message contient un champ nommé b, le nom du champ analysé devient value_b.

          Remarque

          La valeur de value.fields-prefix ne peut pas être un préfixe de la valeur de key.fields-prefix.

          metadata.list

          Colonnes de métadonnées transmises au puits en aval.

          Non

          String

          Les colonnes de métadonnées disponibles incluent topic, partition, offset, timestamp, timestamp-type, headers et leader-epoch. Séparez les noms de colonnes par des virgules.

          scan.value.initial-schemas.ddls

          Instructions DDL définissant le schéma initial pour des tables spécifiques.

          Non

          String

          Utilisez un point-virgule (;) pour séparer plusieurs instructions DDL. Par exemple, utilisez CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT); pour spécifier respectivement le schéma initial des tables db1.t1 et db1.t2.

          Le schéma de table défini dans l'instruction DDL doit être cohérent avec la table de destination du puits et respecter la syntaxe Flink SQL.

          Remarque

          Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 11,5.

          ingestion.ignore-errors

          Ignore les erreurs d'analyse des données.

          Non

          Boolean

          false

          Remarque

          Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 11,5.

          ingestion.error-tolerance.max-count

          Nombre maximal d'erreurs d'analyse tolérées avant l'échec du job. Ne prend effet que si ingestion.ignore-errors est défini sur true.

          Non

          Integer

          -1

          Ce paramètre s'applique uniquement lorsque ingestion.ignore-errors est défini sur true. Une valeur de -1 indique une tolérance illimitée, ce qui signifie que les exceptions d'analyse n'entraîneront pas l'échec du job.

          Remarque

          Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 11,5.

          scan.duplicate-field.strategy

          Spécifie comment gérer les noms de champs en double analysés à partir des parties clé et valeur.

          Non

          String

          EXCEPTION

          Valeurs valides :

          • EXCEPTION : Lève une exception lorsque des champs en double existent dans la clé et la valeur. Il s'agit du comportement par défaut dans VVR 11.6 et versions antérieures.

          • PREFER_KEY : Utilise la valeur du champ clé en cas de duplication des champs.

          • PREFER_VALUE : Utilise la valeur du champ valeur en cas de duplication des champs.

          Remarque

          Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 11,7.

          • Paramètres du format Debezium JSON

            Paramètre

            Obligatoire

            Type

            Valeur par défaut

            Description

            debezium-json.distributed-tables

            Non

            Boolean

            false

            Définissez sur true si les données d'une seule table Debezium JSON sont distribuées sur plusieurs partitions.

            Remarque

            Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 8.0.11.

            Important

            La modification de ce paramètre nécessite un démarrage sans état.

            debezium-json.schema-include

            Non

            Boolean

            false

            Inclut un schéma dans le message Debezium JSON. Cela correspond à la propriété value.converter.schemas.enable dans la configuration Debezium Kafka Connect.

            Valeurs valides :

            • true : Le message Debezium JSON contient un schéma.

            • false : Le message Debezium JSON ne contient pas de schéma.

            debezium-json.ignore-parse-errors

            Non

            Boolean

            false

            Valeurs valides :

            • true : Ignore les lignes provoquant une exception d'analyse.

            • false : Lève une erreur et le job échoue.

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

            Non

            Boolean

            false

            Analyse tous les types primitifs en tant que String lors de l'analyse du schéma de la table.

            Valeurs valides :

            • true : Analyse tous les types primitifs en tant que String.

            • false : Analyse les types selon les règles par défaut.

            debezium-json.infer-schema.string-type-inference

            Non

            Boolean

            true

            Indique s'il faut déduire les champs de chaîne comme TIME, DATE ou TIMESTAMP. Si vous définissez ce paramètre sur false, le connecteur ignore cette inférence et conserve les champs en tant que STRING.

            Remarque

            Ce paramètre n'est pris en charge qu'à partir de Ververica Runtime (VVR) 11,8.

          • Paramètres du format Canal JSON

            Parameter

            Required

            Type

            Default

            Description

            canal-json.distributed-tables

            No

            Boolean

            false

            Si les données d'une seule table au format Canal JSON sont réparties sur plusieurs partitions, vous devez activer cette option.

            Remarque

            Cette option de configuration est prise en charge uniquement dans Ververica Runtime (VVR) 8.0.11 et versions ultérieures.

            Important

            La modification de ce paramètre nécessite un démarrage sans état.

            canal-json.database.include

            No

            String

            Expression régulière facultative permettant de filtrer les journaux de modifications (changelogs) selon le champ de métadonnées database présent dans les enregistrements Canal. Seuls les enregistrements provenant des bases de données correspondantes sont traités. L'expression régulière est compatible avec la classe Pattern de Java.

            canal-json.table.include

            No

            String

            Expression régulière facultative permettant de filtrer les journaux de modifications (changelogs) selon le champ de métadonnées table présent dans les enregistrements Canal. Seuls les enregistrements provenant des tables correspondantes sont traités. L'expression régulière est compatible avec la classe Pattern de Java.

            canal-json.ignore-parse-errors

            No

            Boolean

            false

            Valeurs valides :

            • true : ignore la ligne actuelle si une exception d'analyse se produit.

            • false : génère une erreur et empêche le démarrage du job.

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

            No

            Boolean

            false

            Analyse tous les types primitifs sous forme de String lors de l'analyse du schéma de la table.

            Valeurs valides :

            • true : analyse tous les types primitifs sous forme de String.

            • false : analyse les types selon les règles par défaut.

            canal-json.infer-schema.string-type-inference

            No

            Boolean

            true

            Indique s'il faut déduire les champs de type chaîne comme étant de type TIME, DATE ou TIMESTAMP. Si vous définissez ce paramètre sur false, le connecteur ignore cette inférence et conserve les champs en tant que STRING.

            Remarque

            Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.8 et versions ultérieures.

            canal-json.infer-schema.strategy

            No

            String

            AUTO

            Stratégie d'analyse du schéma de la table.

            Valeurs valides :

            • AUTO : analyse automatiquement le schéma à partir des données JSON. Cette option est recommandée si les données ne contiennent pas de champ sqlType, afin d'éviter les échecs d'analyse.

            • SQL_TYPE : analyse le schéma à partir du tableau sqlType présent dans les données Canal JSON. Nous vous recommandons de définir cette option sur SQL_TYPE pour obtenir des types plus précis si les données contiennent un champ sqlType.

            • MYSQL_TYPE : analyse le schéma à partir du tableau mysqlType présent dans les données Canal JSON.

            Si les données Canal JSON dans Kafka contiennent le champ sqlType et que vous avez besoin d'un mappage de types plus précis, définissez canal-json.infer-schema.strategy sur SQL_TYPE.

            Pour plus d'informations sur les règles de mappage des types sqlType, consultez la section Analyse du schéma Canal JSON.

            Remarque
            • Cette configuration est prise en charge uniquement dans Ververica Runtime (VVR) 11.1 et versions ultérieures.

            • La valeur MYSQL_TYPE est prise en charge dans Ververica Runtime (VVR) 11.3 et versions ultérieures.

            canal-json.mysql.treat-mysql-timestamp-as-datetime-enabled

            No

            Boolean

            true

            Mappe le type MySQL TIMESTAMP au type CDC TIMESTAMP.

            • true : le type MySQL TIMESTAMP est mappé au type CDC TIMESTAMP.

            • false : le type MySQL TIMESTAMP est mappé au type CDC TIMESTAMP_LTZ.

            canal-json.mysql.treat-tinyint1-as-boolean.enabled

            No

            Boolean

            true

            Lorsque vous utilisez la stratégie d'analyse MYSQL_TYPE, ce paramètre contrôle si le type MySQL TINYINT(1) doit être mappé au type CDC BOOLEAN.

            • true : le type MySQL TINYINT(1) est mappé au type CDC BOOLEAN.

            • false : le type MySQL TINYINT(1) est mappé au type CDC TINYINT(1).

            Cette option s'applique uniquement lorsque canal-json.infer-schema.strategy est défini sur MYSQL_TYPE.

          • Paramètres de format JSON

            Parameter

            Required

            Type

            Default

            Description

            json.timestamp-format.standard

            No

            String

            SQL

            Format d'horodatage pour les données d'entrée et de sortie.

            • SQL : analyse les horodatages d'entrée au format yyyy-MM-dd HH:mm:ss.s{precision}, par exemple 2020-12-30 12:13:14.123.

            • ISO-8601 : analyse les horodatages d'entrée au format yyyy-MM-ddTHH:mm:ss.s{precision}, par exemple 2020-12-30T12:13:14.123.

            json.ignore-parse-errors

            No

            Boolean

            false

            Valeurs valides :

            • true : ignore la ligne actuelle si une exception d'analyse se produit.

            • false : génère une erreur et empêche le démarrage du job.

            json.infer-schema.primitive-as-string

            No

            Boolean

            false

            Analyse tous les types primitifs sous forme de String lors de l'analyse du schéma de la table.

            Valeurs valides :

            • true : analyse tous les types primitifs sous forme de String.

            • false : analyse les types selon les règles par défaut.

            json.infer-schema.string-type-inference

            No

            Boolean

            true

            Indique s'il faut déduire les champs de type chaîne comme étant de type TIME, DATE ou TIMESTAMP. Si vous définissez ce paramètre sur false, le connecteur ignore cette inférence et conserve les champs en tant que STRING.

            Remarque

            Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.8 et versions ultérieures.

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

            No

            Boolean

            false

            Développe récursivement les colonnes imbriquées dans les données JSON. Valeurs valides :

            • true : développe récursivement les colonnes imbriquées.

            • false : traite les colonnes imbriquées comme des String.

            json.decode.parser-table-id.fields

            No

            String

            Utilise les valeurs des champs JSON spécifiés pour générer un tableId lors de l'analyse des données au format JSON. Les valeurs de plusieurs champs sont concaténées par une virgule anglaise ,. Par exemple, si les données JSON sont {"col0":"a", "col1","b", "col2","c"}, le résultat généré est le suivant :

            Configuration

            tableId

            col0

            a

            col0,col1

            a.b

            col0,col1,col2

            a.b.c

            json.infer-schema.fixed-types

            No

            String

            Lors de l'analyse des données JSON, vous pouvez spécifier les types de données pour certains champs. Utilisez une virgule , pour séparer plusieurs champs. Par exemple, id BIGINT, name VARCHAR(10) spécifie que le champ id est de type BIGINT et que le champ name est de type VARCHAR(10).

            Remarque
            • Cette option de configuration est prise en charge uniquement dans Ververica Runtime (VVR) 11.5 et versions ultérieures.

            • Lorsque vous utilisez cette configuration avec Ververica Runtime (VVR) version 11.5, vous devez également ajouter la configuration scan.max.pre.fetch.records: 0.

            json.decode.empty-value-as-delete.enabled

            No

            Boolean

            false

            Indique s'il faut analyser les messages tombstone (avec une valeur vide) dans un topic Kafka compacté comme des événements DELETE. Cette option est utilisée dans les scénarios où une valeur vide représente une sémantique de suppression, tels que la mise en miroir de topics compactés ou les signaux de suppression CDC.

            Remarque

            Cette option de configuration est prise en charge uniquement dans Ververica Runtime (VVR) 11.7 et versions ultérieures.

        • Paramètres de la table de destination (Sink)

          Parameter

          Description

          Required

          Type

          Default

          Remarks

          type

          Type de destination (Sink).

          Yes

          String

          La valeur doit être kafka.

          name

          Nom de la destination (Sink).

          No

          String

          Aucun

          topic

          Nom du topic Kafka.

          No

          String

          Si ce paramètre est spécifié, toutes les données sont écrites dans ce topic.

          Remarque

          Si ce paramètre n'est pas spécifié, chaque enregistrement est écrit dans un topic nommé d'après son TableID. Le TableID est construit en joignant les noms de la base de données et de la table par un point (.), par exemple databaseName.tableName.

          partition.strategy

          Stratégie d'écriture des partitions Kafka.

          No

          String

          all-to-zero

          Valeurs valides :

          • all-to-zero (par défaut) : écrit toutes les données dans la partition 0.

          • hash-by-key : écrit les données dans les partitions en fonction de la valeur de hachage de la clé primaire. Cela garantit que les enregistrements ayant la même clé primaire sont écrits dans la même partition, préservant ainsi leur ordre.

          sink.tableId-to-topic.mapping

          Mappage des noms de tables amont vers les noms de topics Kafka aval.

          No

          String

          Séparez les mappages par des points-virgules (;). Dans chaque mappage, séparez le nom de la table amont et le nom du topic Kafka aval par deux-points (:). Vous pouvez utiliser une expression régulière pour le nom de la table. Pour mapper plusieurs tables vers le même topic, séparez les noms de table par des virgules (,). Par exemple : mydb.mytable1:topic1;mydb.mytable2:topic2.

          Remarque

          Ce paramètre vous permet de modifier le topic mappé tout en conservant les informations relatives au nom de la table d'origine.

          • Paramètres de format Canal JSON

            Parameter

            Required

            Type

            Default

            Description

            canal-json.serialize.update.keep-changed-fields-only

            No

            Boolean

            false

            Indique si la section old d'un message UPDATE au format Canal JSON contient uniquement les anciennes valeurs des champs qui ont changé.

            Remarque

            Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.8 et versions ultérieures.

          • Paramètres de format Debezium JSON

            Parameter

            Required

            Type

            Default

            Description

            debezium-json.include-schema.enabled

            No

            Boolean

            false

            Inclut les informations de schéma dans les données Debezium JSON.

            debezium-json.emit.full-table-id.enabled

            No

            Boolean

            false

            Écrit l'identifiant complet de la table en trois parties dans les champs de métadonnées Debezium JSON.

            Si ce paramètre est activé, le mappage est le suivant :

            CDC Table ID Part

            Debezium JSON Key

            Namespace

            db

            Schema

            schema

            Table

            table

            Si ce paramètre est désactivé, le mappage est le suivant :

            CDC Table ID Part

            Debezium JSON Key

            Namespace

            Non mappé

            Schema

            db

            Table

            table

            Remarque

            Ce paramètre est pris en charge uniquement dans Ververica Runtime (VVR) 11.6 et versions ultérieures.

        Exemples

        • Utiliser Kafka comme source Flink CDC :

          source:
            type: kafka
            name: Kafka source
            properties.bootstrap.servers: ${kafka.bootstraps.server}
            topic: ${kafka.topic}
            value.format: ${value.format}
            scan.startup.mode: ${scan.startup.mode}
           
          sink:
            type: hologres
            name: Hologres sink
            endpoint: <yourEndpoint>
            dbname: <yourDbname>
            username: ${secret_values.ak_id}
            password: ${secret_values.ak_secret}
            sink.type-normalize-strategy: BROADEN
        • Utiliser Kafka comme destination (Sink) Flink CDC :

          source:
            type: mysql
            name: MySQL Source
            hostname: ${secret_values.mysql.hostname}
            port: ${mysql.port}
            username: ${secret_values.mysql.username}
            password: ${secret_values.mysql.password}
            tables: ${mysql.source.table}
            server-id: 8601-8604
          
          sink:
            type: kafka
            name: Kafka Sink
            properties.bootstrap.servers: ${kafka.bootstraps.server}
          
          route:
            - source-table: ${mysql.source.table}
              sink-table: ${kafka.topic}

          Le module route spécifie le topic Kafka de destination pour la table source.

        Remarque

        Par défaut, la fonctionnalité de création automatique de topics est désactivée pour ApsaraMQ for Kafka. Pour plus d'informations, consultez la section FAQ sur la création automatique de topics. Vous devez créer le topic avant d'écrire des données dans ApsaraMQ for Kafka. Pour plus d'informations, consultez la section Étape 3 : Créer des ressources.

        Lorsque la fonction de checkpointing est activée, la source Kafka valide l'offset actuel du consommateur auprès de Kafka à la fin d'un checkpoint. Cela garantit que l'état du checkpoint Flink est cohérent avec l'offset validé sur le broker Kafka. Si le checkpointing est désactivé, la source Kafka s'appuie sur le mécanisme interne de validation périodique automatique des offsets du consommateur Kafka. Cette fonctionnalité est contrôlée par les propriétés du consommateur Kafka enable.auto.commit et auto.commit.interval.ms.

        Remarque

        La source Kafka ne s'appuie pas sur les offsets validés pour la tolérance aux pannes et la récupération. La validation des offsets sert uniquement à surveiller la progression du consommateur Kafka et du groupe de consommateurs.

      • Autres propriétés

        Outre les propriétés mentionnées, vous pouvez utiliser setProperties(Properties) et setProperty(String, String) pour définir n'importe quelle Property du Kafka Source et de son consommateur Kafka sous-jacent. Le Kafka Source propose les propriétés spécifiques suivantes.

        Parameter

        Description

        client.id.prefix

        Préfixe de l'ID client pour le consommateur Kafka.

        partition.discovery.interval.ms

        Intervalle de découverte des partitions en millisecondes. La valeur -1 désactive la découverte dynamique des partitions.

        Remarque

        En Batch Mode, cette propriété est automatiquement définie sur -1.

        register.consumer.metrics

        Enregistre les métriques du consommateur Kafka dans Flink.

        Autres configurations du consommateur Kafka

        Pour obtenir la liste complète des configurations du consommateur Kafka, consultez la documentation officielle d'Apache Kafka.

        Important

        Pour garantir un fonctionnement correct, le connecteur DataStream Connector de Kafka écrase les propriétés configurées manuellement suivantes :

        • key.deserializer est toujours remplacé par org.apache.kafka.common.serialization.ByteArrayDeserializer.

        • value.deserializer est toujours remplacé par org.apache.kafka.common.serialization.ByteArrayDeserializer.

        • auto.offset.reset.strategy est remplacé par la stratégie fournie par OffsetsInitializer.

        L'exemple suivant montre comment configurer un consommateur Kafka pour utiliser le mécanisme SASL PLAIN et fournir une configuration JAAS.

        KafkaSource.builder()
            .setProperty("sasl.mechanism", "PLAIN")
            .setProperty("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"username\" password=\"password\";")
      • Monitoring

        Le Kafka Source expose des métriques via le système de métriques de Flink à des fins de surveillance et de diagnostic.

        • Portée des métriques

          Toutes les métriques du lecteur source Kafka sont enregistrées dans le groupe de métriques KafkaSourceReader, qui est un sous-groupe du groupe de métriques de l'opérateur. Les métriques relatives à une partition de topic spécifique sont enregistrées dans le sous-groupe KafkaSourceReader.topic.<topic_name>.partition.<partition_id>.

          Par exemple, la métrique d'Offset actuel du consommateur (currentOffset) pour la partition 1 du topic « my-topic » est disponible à l'emplacement .operator.KafkaSourceReader.topic.my-topic.partition.1.currentOffset. Le nombre de validations réussies (commitsSucceeded) est disponible à l'emplacement .operator.KafkaSourceReader.commitsSucceeded.

        • Liste des métriques

          Metric

          Description

          Scope

          currentOffset

          L'Offset actuel du consommateur pour une partition.

          TopicPartition

          committedOffset

          Dernier Offset validé pour une partition.

          TopicPartition

          commitsSucceeded

          Nombre total de validations d'offset réussies.

          KafkaSourceReader

          commitsFailed

          Nombre de validations échouées

          KafkaSourceReader

        • Métriques du consommateur Kafka

          Les métriques du consommateur Kafka sous-jacent sont enregistrées dans le groupe de métriques KafkaSourceReader.KafkaConsumer. Par exemple, la métrique records-consumed-total est enregistrée à l'emplacement .operator.KafkaSourceReader.KafkaConsumer.records-consumed-total.

          Utilisez la propriété register.consumer.metrics pour indiquer si les métriques du consommateur Kafka doivent être enregistrées. Cette option est activée par défaut (true). Pour plus d'informations sur les métriques du consommateur Kafka, consultez la documentation Apache Kafka.

      • Créer un Kafka sink

        Le Kafka Sink de Flink écrit un flux de données dans un ou plusieurs topics Kafka.

        DataStream<String> stream = ...
        
        Properties kafkaProperties = new Properties();
        kafkaProperties.setProperty("bootstrap.servers", "localhost:9092");
        
        KafkaSink<String> sink = KafkaSink.<String>builder()
                .setKafkaProducerConfig(kafkaProperties)
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.builder()
                                .setTopic("my-topic")
                                .setValueSerializationSchema(new SimpleStringSchema())
                                .build())
                .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
                .build();
        
        stream.sinkTo(sink);

        Pour créer un Kafka Sink, vous devez configurer les propriétés suivantes.

        Parameter

        Description

        Propriétés du client Kafka

        La Property bootstrap.servers est requise. Elle spécifie une liste de brokers Kafka séparés par des virgules.

        Sérialiseur d'enregistrements

        Vous devez fournir un KafkaRecordSerializationSchema pour convertir les données d'entrée en un ProducerRecord Kafka. Flink fournit un générateur de schéma qui offre des composants courants, tels que la sérialisation des clés et des valeurs des messages, la sélection des topics et le partitionnement des messages. Vous pouvez également implémenter les interfaces correspondantes pour un contrôle plus granulaire. La méthode ProducerRecord<byte[], byte[]> serialize(T element, KafkaSinkContext context, Long timestamp) est appelée pour chaque enregistrement entrant afin de générer un ProducerRecord à écrire dans Kafka.

        Le ProducerRecord offre un contrôle précis sur la manière dont chaque enregistrement est écrit dans Kafka, ce qui vous permet de :

        • Définir le Topic de destination.

        • Définir la Key du Message.

        • Spécifier la Partition de destination.

        Garantie de livraison

        Le paramètre bootstrap.servers est requis et spécifie une liste de brokers Kafka séparés par des virgules.

        Garantie de livraison

        Lorsque les points de contrôle Flink sont activés, le Kafka Sink de Flink peut fournir une sémantique exactement une fois. En plus d'activer les points de contrôle, vous pouvez utiliser le paramètre DeliveryGuarantee pour spécifier différentes garanties de livraison. Le paramètre DeliveryGuarantee propose les options suivantes :

        • DeliveryGuarantee.NONE : (par défaut) Flink ne fournit aucune garantie. Des données peuvent être perdues ou dupliquées.

        • DeliveryGuarantee.AT_LEAST_ONCE : garantit qu'aucune donnée n'est perdue, mais des duplications peuvent se produire.

        • DeliveryGuarantee.EXACTLY_ONCE : utilise les transactions Kafka pour fournir une sémantique exactement une fois.

          Remarque

          Lors de l'utilisation de la sémantique EXACTLY_ONCE, consultez les Considérations relatives à la sémantique exactement une fois.

      • Flink CDC

        Utilisez le connecteur Kafka en tant que source ou sink pour créer des tâches YAML pour Flink CDC.

        Limitations

        • Utilisez Realtime Compute for Apache Flink (VVR) version 11.1 ou ultérieure pour ingérer des données Flink CDC à partir d'une source de données Kafka.

        • Seuls les formats JSON, Debezium JSON et Canal JSON sont pris en charge.

        • Seule la version 8.0.11 ou ultérieure de Realtime Compute for Apache Flink (VVR) prend en charge la lecture des données d'une seule table distribuée sur plusieurs partitions.

        Syntaxe

        source:
          type: kafka
          name: Kafka source
          properties.bootstrap.servers: localhost:9092
          topic: ${kafka.topic}
        sink:
          type: kafka
          name: Kafka Sink
          properties.bootstrap.servers: localhost:9092

        Stratégies d'analyse et d'évolution du schéma

        Le connecteur Kafka conserve les schémas de toutes les tables actuellement connues.

        Initialisation du schéma de table

        Un schéma de table comprend les colonnes et les types de données, les noms de base de données et de table, ainsi que les clés primaires. Les sections suivantes décrivent comment initialiser chacun de ces éléments.

        • Informations sur les colonnes et les types de données

        Une tâche Flink CDC peut déduire automatiquement les colonnes et les types de données à partir des données, mais il peut être nécessaire de les définir explicitement pour certaines tables. Il existe trois stratégies d'initialisation de schéma, selon le niveau de contrôle souhaité sur les types :

        1. Inférence automatique complète du schéma

        Avant de lire les données depuis Kafka, le connecteur Kafka tente de consommer jusqu'à scan.max.pre.fetch.records messages de chaque partition, analyse le schéma de chaque message et fusionne ces schémas pour initialiser le schéma de la table. Un événement de création de table est ensuite généré sur la base de ce schéma initialisé avant que les données ne soient réellement consommées.

        Remarque

        Pour les formats Debezium JSON et Canal JSON, les informations de table sont contenues dans chaque message. Les messages pré-extraits sur la base du paramètre scan.max.pre.fetch.records peuvent contenir des données provenant de plusieurs tables. Par conséquent, le nombre d'enregistrements pré-extraits pour une seule table ne peut pas être déterminé. La pré-extraction et l'initialisation du schéma sont effectuées une seule fois pour chaque partition avant que ses messages ne soient consommés et traités. Si des données pour une nouvelle table apparaissent ultérieurement, le schéma analysé à partir du premier enregistrement de cette table est utilisé comme schéma initial, et le schéma n'est ni pré-extrait ni initialisé à nouveau.

        Important

        La distribution des données d'une seule table sur plusieurs partitions est prise en charge uniquement dans Ververica Runtime (VVR) version 8.0.11 et ultérieure, et nécessite de définir l'option de configuration debezium-json.distributed-tables ou canal-json.distributed-tables sur true.

        1. Spécification d'un schéma de table initial

        Dans certains cas, il peut être nécessaire de définir explicitement le schéma de table initial, par exemple lors de l'écriture de données depuis Kafka vers une table aval préexistante. Dans ce cas, vous pouvez le faire en ajoutant le paramètre scan.value.initial-schemas.ddls. Voici un exemple de configuration :

        source:
          type: kafka
          name: Kafka Source
          properties.bootstrap.servers: host:9092
          topic: test-topic
          value.format: json
          scan.startup.mode: earliest-offset
          # Set the initial table schema
          scan.value.initial-schemas.ddls: CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);

        L'instruction DDL doit correspondre au schéma de la table cible. Cette configuration spécifie le type initial de la colonne id comme BIGINT et de la colonne name comme VARCHAR(10) pour la table db1.t1, et le type initial de la colonne id comme BIGINT pour la table db1.t2.

        Les instructions DDL utilisent la syntaxe Flink SQL.

        1. Définition de types fixes pour des champs spécifiques

        Il peut être nécessaire de verrouiller certains champs sur un type de données fixe. Par exemple, des champs qui seraient normalement inférés comme TIMESTAMP pourraient devoir être produits sous forme de chaînes. Dans ce cas, vous pouvez ajouter le paramètre json.infer-schema.fixed-types pour spécifier le schéma de table initial. Ce paramètre n'est valide que lorsque le format du message est JSON. Voici un exemple de configuration :

        source:
          type: kafka
          name: Kafka Source
          properties.bootstrap.servers: host:9092
          topic: test-topic
          value.format: json
          scan.startup.mode: earliest-offset
          # Set specific fields to a fixed type
          json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
          scan.max.pre.fetch.records: 0

        Cette configuration spécifie que tous les champs id sont de type BIGINT et que tous les champs name sont de type VARCHAR(10).

        Les types de données sont cohérents avec les types Flink SQL.

        • Informations sur la base de données et la table

          • Pour les formats Canal JSON et Debezium JSON, le connecteur analyse les informations de table, y compris le nom de la base de données et de la table, à partir de chaque message.

          • Pour le format JSON, par défaut, les informations de table ne contiennent que le nom de la table, qui correspond au nom du topic contenant les données. Si vos données contiennent des informations sur la base de données et la table, vous pouvez utiliser le paramètre json.infer-schema.fixed-types pour spécifier les champs contenant ces informations. Ces champs sont ensuite mappés aux noms de base de données et de table. Voici un exemple de configuration :

            source:
              type: kafka
              name: Kafka Source
              properties.bootstrap.servers: host:9092
              topic: test-topic
              value.format: json
              scan.startup.mode: earliest-offset
              # Use the value of the col1 field as the database name and the value of the col2 field as the table name
              json.decode.parser-table-id.fields: col1,col2

            Avec cette configuration, le connecteur envoie chaque enregistrement à une table où le nom de la base de données est la valeur du champ col1 et le nom de la table est la valeur du champ col2.

        • Informations sur la clé primaire

          • Pour le format Canal JSON, le champ pkNames dans les données JSON définit la clé primaire de la table.

          • Pour les formats Debezium JSON et JSON, les données ne contiennent pas d'informations sur la clé primaire. Vous pouvez ajouter manuellement des clés primaires aux tables en utilisant des règles transform :

            transform:
              - source-table: \.*.\.*
                projection: \*
                primary-keys: key1, key2

        Analyse du schéma et évolution du schéma

        Après l'initialisation du schéma de table, si schema.inference.strategy est défini sur static, le connecteur Kafka analyse la valeur du message de chaque message en fonction du schéma de table initial et ne génère pas d'événements de modification de schéma. Si schema.inference.strategy est défini sur continuous, le connecteur Kafka analyse la valeur du message de chaque message Kafka, identifie ses colonnes physiques et compare le schéma résultant avec le schéma actuellement maintenu. Si les schémas sont incohérents, le connecteur tente de les fusionner et génère un événement de modification de schéma de table correspondant. Les règles de fusion sont les suivantes :

        • Si les colonnes physiques analysées contiennent des champs absents du schéma actuel, ces champs sont ajoutés au schéma, et un événement est généré pour les ajouter en tant que colonnes acceptant les valeurs nulles.

        • Si les colonnes physiques analysées ne contiennent pas de champs existant dans le schéma actuel, ces champs sont conservés et leurs valeurs sont renseignées avec NULL. Aucun événement de suppression de colonne n'est généré.

        • Les colonnes portant le même nom sont gérées comme suit :

          • Si les colonnes ont le même type de données mais une précision différente, le type avec la précision la plus élevée est utilisé, et un événement de modification de type de colonne est généré.

          • Si les colonnes ont des types de données différents, le système trouve le plus petit type parent commun dans l'arborescence hiérarchique des types ci-dessous. Le système utilise ensuite ce type parent commun pour la colonne et génère un événement de modification de type de colonne.

            image

        • Stratégies d'évolution de schéma prises en charge :

          • Ajout d'une colonne : le connecteur ajoute la nouvelle colonne à la fin du schéma et synchronise ses données. La nouvelle colonne est définie comme acceptant les valeurs nulles.

          • Suppression d'une colonne : aucun événement de suppression de colonne n'est généré. À la place, les données ultérieures pour cette colonne sont renseignées avec NULL.

          • Renommage d'une colonne : le connecteur traite cette opération comme une suppression de l'ancienne colonne et un ajout d'une nouvelle. La nouvelle colonne est ajoutée à la fin du schéma, et les valeurs de la colonne d'origine sont renseignées avec NULL.

          • Modification du type d'une colonne :

            • Pour les sinks aval qui prennent en charge les modifications de type de colonne, une tâche Flink CDC peut gérer les changements de type (par exemple, de INT vers BIGINT) si le sink aval est configuré pour les traiter. Cette capacité dépend des règles de modification de type de colonne prises en charge par le sink spécifique. Reportez-vous à la documentation de votre sink pour connaître les règles prises en charge.

            • Pour les sinks aval qui ne prennent pas en charge les modifications de type de colonne, tels que Hologres, vous pouvez utiliser l'élargissement de type. Cette fonctionnalité crée une table avec des types de données plus larges dans le sink aval au démarrage de la tâche. Lorsqu'un type de colonne change, le système peut tolérer le changement tant que le nouveau type tient dans le type plus large défini dans le sink aval.

        • Modifications de schéma non prises en charge :

          • Modifications des contraintes, telles que les clés primaires ou les index.

          • Changement d'une colonne de NOT NULL vers NULLABLE.

        • Analyse du schéma Canal JSON

          Les données Canal JSON peuvent contenir un champ facultatif sqlType, qui enregistre des informations de type précises pour les colonnes de données. Pour obtenir un schéma plus précis, vous pouvez définir canal-json.infer-schema.strategy sur SQL_TYPE pour utiliser les types du champ sqlType. Les mappages de types sont les suivants :

          Type JDBC

          Code de type

          Type CDC

          BIT

          -7

          BOOLEAN

          BOOLEAN

          16

          TINYINT

          -6

          TINYINT

          SMALLINT

          5

          SMALLINT

          INTEGER

          4

          INT

          BIGINT

          -5

          BIGINT

          DECIMAL

          3

          DECIMAL(38,18)

          NUMERIC

          2

          REAL

          7

          FLOAT

          FLOAT

          6

          DOUBLE

          8

          DOUBLE

          BINARY

          -2

          BYTES

          VARBINARY

          -3

          LONGVARBINARY

          -4

          BLOB

          2004

          DATE

          91

          DATE

          TIME

          92

          TIME

          TIMESTAMP

          93

          TIMESTAMP

          CHAR

          1

          STRING

          VARCHAR

          12

          LONGVARCHAR

          -1

          Autres types de données

        Tolérance et collecte des données incorrectes

        Votre source de données Kafka peut contenir des enregistrements mal formés, communément appelés données incorrectes. Pour éviter que votre tâche n'échoue et ne redémarre de manière répétée, vous pouvez la configurer pour ignorer ces enregistrements invalides. Par exemple :

        source:
          type: kafka
          name: Kafka Source
          properties.bootstrap.servers: host:9092
          topic: test-topic
          value.format: json
          scan.startup.mode: earliest-offset
          # Enable Dirty Data Tolerance
          ingestion.ignore-errors: true
          # Tolerate up to 1000 dirty data records
          ingestion.error-tolerance.max-count: 1000

        Avec cette configuration, la tâche continue de s'exécuter tant qu'elle ne rencontre pas plus de 1 000 enregistrements incorrects. Une fois ce seuil dépassé, la tâche échoue afin que vous puissiez examiner vos données.

        Pour garantir que votre tâche n'échoue jamais en raison de données incorrectes, utilisez la configuration suivante :

        source:
          type: kafka
          name: Kafka Source
          properties.bootstrap.servers: host:9092
          topic: test-topic
          value.format: json
          scan.startup.mode: earliest-offset
          # Enable Dirty Data Tolerance
          ingestion.ignore-errors: true
          # Tolerate all dirty data records
          ingestion.error-tolerance.max-count: -1

        Bien que la tolérance aux données incorrectes permette à votre tâche de continuer à s'exécuter, vous souhaiterez peut-être également inspecter les enregistrements problématiques. Vous pouvez aussi analyser les données incorrectes pour améliorer vos producteurs Kafka. Comme décrit dans la section Collecte des données incorrectes, vous pouvez afficher les données incorrectes de la tâche dans les journaux TaskManager. Par exemple :

        source:
          type: kafka
          name: Kafka Source
          properties.bootstrap.servers: host:9092
          topic: test-topic
          value.format: json
          scan.startup.mode: earliest-offset
          # Enable Dirty Data Tolerance
          ingestion.ignore-errors: true
          # Tolerate all dirty data records
          ingestion.error-tolerance.max-count: -1
        
        pipeline:
          dirty-data.collector:
            # Write dirty data to the TaskManager log file
            type: logger

        Mappage des noms de table et des topics

        Lorsque Kafka sert de sink Flink CDC, le format du message (tel que Debezium JSON ou Canal JSON) intègre le nom de table d'origine. Les consommateurs aval utilisent généralement ce nom intégré comme identifiant de table plutôt que le nom du topic, il est donc important de configurer correctement le mappage entre les noms de table et les topics.

        Supposons que vous deviez synchroniser deux tables d'une base de données MySQL : mydb.mytable1 et mydb.mytable2. Les stratégies de mappage suivantes sont disponibles :

        1. Aucune stratégie de mappage

        Sans aucune stratégie de mappage, les données de chaque table sont écrites dans un topic nommé au format <Nom de la base de données>.<Nom de la table>. Par conséquent, les données de mydb.mytable1 sont écrites dans un topic nommé mydb.mytable1, et les données de mydb.mytable2 sont écrites dans un topic nommé mydb.mytable2. Voici un exemple de configuration :

        source:
          type: mysql
          name: MySQL Source
          hostname: ${secret_values.mysql.hostname}
          port: ${mysql.port}
          username: ${secret_values.mysql.username}
          password: ${secret_values.mysql.password}
          tables: mydb.mytable1,mydb.mytable2
          server-id: 8601-8604
        
        sink:
          type: kafka
          name: Kafka Sink
          properties.bootstrap.servers: ${kafka.bootstraps.server}

        2. Mappage par règle de routage (non recommandé)

        Vous souhaitez peut-être écrire des données dans un topic spécifique au lieu d'utiliser le format par défaut <Nom de la base de données>.<Nom de la table>. Pour ce faire, vous pouvez configurer une règle de routage. Voici un exemple de configuration :

        source:
          type: mysql
          name: MySQL Source
          hostname: ${secret_values.mysql.hostname}
          port: ${mysql.port}
          username: ${secret_values.mysql.username}
          password: ${secret_values.mysql.password}
          tables: mydb.mytable1,mydb.mytable2
          server-id: 8601-8604
        
        sink:
          type: kafka
          name: Kafka Sink
          properties.bootstrap.servers: ${kafka.bootstraps.server}
          
         route:
          - source-table: mydb.mytable1,mydb.mytable2
            sink-table: mytable

        Dans ce cas, toutes les données de mydb.mytable1 et mydb.mytable2 sont écrites dans un seul topic nommé mytable.

        Cependant, une règle de routage qui modifie le topic de destination change également le nom de la table dans le message Kafka (au format Debezium JSON ou Canal JSON). Le nom de la table dans tous les messages Kafka devient mytable. Cela peut entraîner un comportement inattendu dans les systèmes qui consomment des messages de ce topic.

        3. Mappage avec sink.tableId-to-topic.mapping (recommandé)

        Pour mapper les noms de table aux topics tout en conservant le nom de la table source d'origine, utilisez le paramètre sink.tableId-to-topic.mapping. Voici un exemple de configuration :

        source:
          type: mysql
          name: MySQL Source
          hostname: ${secret_values.mysql.hostname}
          port: ${mysql.port}
          username: ${secret_values.mysql.username}
          password: ${secret_values.mysql.password}
          tables: mydb.mytable1,mydb.mytable2
          server-id: 8601-8604
          sink.tableId-to-topic.mapping: mydb.mytable1,mydb.mytable2:mytable
        
        sink:
          type: kafka
          name: Kafka Sink
          properties.bootstrap.servers: ${kafka.bootstraps.server}

        Vous pouvez également utiliser la configuration suivante :

        source:
          type: mysql
          name: MySQL Source
          hostname: ${secret_values.mysql.hostname}
          port: ${mysql.port}
          username: ${secret_values.mysql.username}
          password: ${secret_values.mysql.password}
          tables: mydb.mytable1,mydb.mytable2
          server-id: 8601-8604
          sink.tableId-to-topic.mapping: mydb.mytable1:mytable;mydb.mytable2:mytable
        
        sink:
          type: kafka
          name: Kafka Sink
          properties.bootstrap.servers: ${kafka.bootstraps.server}

        Dans ce cas, toutes les données de mydb.mytable1 et mydb.mytable2 sont écrites dans le topic mytable, et le nom de la table dans les messages Kafka (au format Debezium JSON ou Canal JSON) est conservé en tant que mydb.mytable1 ou mydb.mytable2. Ainsi, les systèmes aval peuvent toujours identifier la table source d'origine de chaque enregistrement.

        Exemples de configuration

        Les exemples suivants montrent des configurations pour des cas d'utilisation courants.

        Lire à partir d'un seul topic

        L'exemple suivant lit le topic customers et écrit les données dans un lac de données Alibaba Cloud :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true

        Pour le format JSON, le nom de table généré est identique au nom du topic par défaut.

        Lire à partir de plusieurs topics

        L'exemple suivant lit les topics dont les noms correspondent à une expression régulière et écrit les données dans StarRocks :

        source:
          type: kafka
          topic-pattern: user_event_.*
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
        
        sink:
          type: starrocks
          jdbc-url: jdbc:mysql://<yourFeHostname>:9030
          load-url: <yourFeHostname>:8030
          username: <yourUsername>
          password: ${secret_values.starrocks_password}
         
          # Optional: For jobs with low data volumes, use a shorter flush interval to prevent data from remaining unwritten for a long time. The default is 300000 milliseconds, or 5 minutes.
          sink.buffer-flush.interval-ms: 5000
          # Optional: If the upstream character set is utf8mb4, set this parameter to 4 to prevent text truncation. The default is 3.
          unicode-char.max-bytes: 4
          # Optional: Specify the number of buckets for automatically created tables. You must explicitly configure this parameter for StarRocks versions earlier than 2.5.7. Later versions can infer the value.
          table.create.num-buckets: 8
          # Optional: Specify the number of replicas for automatically created tables based on your cluster configuration.
          table.create.properties.replication_num: 3
          # Optional: For StarRocks 3.2 and later, enable this feature to accelerate schema changes.
          table.create.properties.fast_schema_evolution: true
          # Note: If a transform changes a primary key, you must also set sink.ignore.update-before to false.
          # Otherwise, the row that uses the old primary key remains in the downstream system.

        Pour le format JSON, le nom de la table générée est identique au nom du topic par défaut.

        Lire les clés et éviter les conflits de champs

        Vous pouvez utiliser l'une des méthodes suivantes pour éviter les erreurs causées par des noms de champs en double dans la clé et la valeur :

        1. Ajoutez des préfixes aux champs pour éviter les conflits :

          source:
            type: kafka
            topic: ${kafka.topic}
            properties.bootstrap.servers: localhost:9092
            properties.group.id: ${kafka.group.id}
            key.format: json
            value.format: json
            # Add the key_ prefix to field names in the key.
            key.fields-prefix: key_
            # Add the value_ prefix to field names in the value.
            value.fields-prefix: value_
            # Optional: Dynamically infer the schema of each record and detect schema changes.
            schema.inference.strategy: continuous
          
          sink:
            type: paimon
            name: Paimon Sink
            catalog.properties.metastore: rest
            catalog.properties.uri: dlf_uri
            catalog.properties.warehouse: your_warehouse
            catalog.properties.token.provider: dlf
            # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
            commit.user: your_job_name
            # Optional: Enable deletion vectors to improve read performance.
            table.properties.deletion-vectors.enabled: true
        2. Configurez une politique de résolution des conflits. Pour plus d'informations, consultez la section scan.duplicate-field.strategy. La configuration suivante privilégie les champs de la clé et ignore les champs portant le même nom dans la valeur :

          source:
            type: kafka
            topic: ${kafka.topic}
            properties.bootstrap.servers: localhost:9092
            properties.group.id: ${kafka.group.id}
            key.format: json
            value.format: json
            # Prefer fields in the key and ignore fields with the same names in the value.
            scan.duplicate-field.strategy: PREFER_KEY
            # Optional: Dynamically infer the schema of each record and detect schema changes.
            schema.inference.strategy: continuous
          
          sink:
            type: paimon
            name: Paimon Sink
            catalog.properties.metastore: rest
            catalog.properties.uri: dlf_uri
            catalog.properties.warehouse: your_warehouse
            catalog.properties.token.provider: dlf
            # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
            commit.user: your_job_name
            # Optional: Enable deletion vectors to improve read performance.
            table.properties.deletion-vectors.enabled: true

        Ajouter des colonnes de métadonnées

        L'exemple suivant lit le topic customers et ajoute les colonnes de métadonnées topic et partition :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
          metadata.list: topic,partition
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true

        Gérer les erreurs d'analyse

        Les erreurs d'analyse entraînent l'échec d'un job. Vous pouvez configurer le job pour qu'il tolère ces erreurs. Cette fonctionnalité est couramment utilisée conjointement avec la collecte de données incorrectes.

        L'exemple suivant ignore toutes les erreurs d'analyse :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
          # Ignore parsing errors. By default, all parsing errors are ignored when this feature is enabled.
          ingestion.ignore-errors: true
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        L'exemple suivant fait échouer le job après 30 erreurs d'analyse :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
          # Ignore parsing errors.
          ingestion.ignore-errors: true
          # Fail the job after 30 parsing errors.
          ingestion.error-tolerance.max-count: 30
          
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        Lire des données JSON

        Les sections suivantes décrivent les méthodes courantes de lecture des données au format JSON.

        Spécifier l'analyse de l'ID de table

        Par défaut, l'ID de table des données JSON correspond au nom du topic. Vous pouvez utiliser les valeurs des champs des données comme ID de table. L'exemple suivant utilise les champs db et tbl comme ID de table :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
          # Use the db and tbl fields as the table ID.
          value.json.decode.parser-table-id.fields: db,tbl
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        Spécifier les types de champs

        Les types de champs sont déduits des valeurs des champs, mais les types inférés peuvent ne pas correspondre à vos attentes. Vous pouvez spécifier des types fixes pour certains champs et ignorer l'inférence et l'évolution ultérieures des types pour ces champs.

        L'exemple suivant fixe les types de quatre champs :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: json
          # Use the db and tbl fields as the table ID.
          value.json.decode.parser-table-id.fields: db,tbl
          # Specify fixed types for selected fields.
          value.json.infer-schema.fixed-types: 'db STRING, tbl STRING, id BIGINT, amount DECIMAL(18, 2)'
          # Continue to dynamically infer fields that are not declared.
          schema.inference.strategy: continuous
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        Lire des données Canal JSON

        Les sections suivantes décrivent les méthodes courantes de lecture des données au format Canal JSON.

        Configurer la politique d'inférence de type

        Par défaut, le connecteur déduit les types de schéma à partir des valeurs des champs lors de la lecture des données Canal JSON.

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: canal-json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        Vous pouvez également déduire le schéma à partir des informations de schéma enregistrées dans les données Canal JSON, telles que les types SQL ou MySQL. L'exemple suivant utilise les types MySQL pour déduire le schéma :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: canal-json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
          # Infer the schema from MySQL type information. You can set this parameter to SQL_TYPE to use SQL type information instead.
          value.canal-json.infer-schema.strategy: MYSQL_TYPE
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true
          
        # Enable a dirty data collector to print records that cannot be parsed.
        pipeline:
          dirty-data.collector:
            name: Logger Dirty Data Collector
            type: logger

        Lire des données Debezium JSON

        L'exemple suivant lit le topic customers et écrit les données dans un lac de données Alibaba Cloud :

        source:
          type: kafka
          topic: customers
          properties.bootstrap.servers: localhost:9092
          properties.group.id: ${kafka.group.id}
          value.format: debezium-json
          # Optional: Dynamically infer the schema of each record and detect schema changes.
          schema.inference.strategy: continuous
        
        sink:
          type: paimon
          name: Paimon Sink
          catalog.properties.metastore: rest
          catalog.properties.uri: dlf_uri
          catalog.properties.warehouse: your_warehouse
          catalog.properties.token.provider: dlf
          # Optional: Specify a commit user. We recommend that you use a different commit user for each job to prevent conflicts.
          commit.user: your_job_name
          # Optional: Enable deletion vectors to improve read performance.
          table.properties.deletion-vectors.enabled: true

        Synchroniser les binlogs MySQL bruts vers Kafka

        L'ingestion de données Flink CDC permet de synchroniser les données binlog MySQL brutes au format Canal JSON. Le job suivant synchronise les binlogs de plusieurs tables vers le topic order_dw_tables :

        source:
          type: mysql
          hostname: #{hostname}
          port: 3306
          username: #{username}
          password: #{password}
          tables: order_dw.\.*
          server-id: 28601-28604
          # Optional: Synchronize data from tables that are created during the incremental phase.
          scan.binlog.newly-added-table.enabled: true
          # Optional: Synchronize table and column comments.
          include-comments.enabled: true
          # Optional: Process unbounded chunks first to prevent potential TaskManager out-of-memory errors.
          scan.incremental.snapshot.unbounded-chunk-first.enabled: true
          # Optional: Enable parsing filters to accelerate reads.
          scan.only.deserialize.captured.tables.changelog.enabled: true
          # Add mysqlType, sqlType, sql, isDdl, and other metadata to the Canal JSON data.
          include-binlog-meta.enable: true
          
        sink:
          type: kafka
          properties.bootstrap.servers: localhost:9092
          topic: order_dw_tables
          # Use the Canal JSON changelog format for Kafka values.
          value.format: canal-json
          # Specify the format used to serialize date and time data.
          value.canal-json.timestamp-format.standard: SQL
          # Write all data to Partition 0 to preserve binlog order.
          partition.strategy: all-to-zero

        Sémantique exactement une fois (Exactly-once)

        • Configurer le niveau d'isolation du consommateur

          Toutes les applications qui consomment des données Kafka doivent définir la propriété isolation.level :

          • read_committed : Lit uniquement les données validées.

          • read_uncommitted (par défaut) : Peut lire des données non validées.

          EXACTLY_ONCE dépend de read_committed. Sinon, les consommateurs peuvent voir des données non validées, ce qui rompt la cohérence.

        • Délai d'expiration des transactions et perte de données

          Lors de la récupération à partir d'un point de contrôle, Realtime Compute for Apache Flink prend uniquement en compte les transactions qui ont été validées avant le début de ce point de contrôle. Si la durée entre l'échec d'un job et son redémarrage dépasse le délai d'expiration de la transaction Kafka, Kafka abandonne automatiquement la transaction ouverte, ce qui peut entraîner une perte de données.

          • La valeur par défaut de transaction.max.timeout.ms pour un broker Kafka est de 15 minutes.

          • Par défaut, Flink Kafka Sink définit le paramètre transaction.timeout.ms sur 1 heure.

          • Vous devez augmenter la valeur de transaction.max.timeout.ms sur le broker pour qu'elle soit supérieure ou égale au paramètre défini dans Flink.

        • Pool de producteurs et points de contrôle simultanés

          Le mode EXACTLY_ONCE utilise un pool de producteurs Kafka de taille fixe. Chaque point de contrôle utilise un producteur issu de ce pool. Si le nombre de points de contrôle simultanés dépasse la taille du pool, le job échoue.

          Configurez la taille du pool de producteurs en fonction du nombre maximal de points de contrôle simultanés.

        • Contraintes de réduction du parallélisme

          Si un job échoue avant la fin du premier point de contrôle, les informations du pool de producteurs d'origine sont perdues lors du redémarrage. Par conséquent, ne réduisez pas le parallélisme du job avant la fin du premier point de contrôle. Si une réduction s'avère nécessaire, le nouveau parallélisme ne doit pas être inférieur à FlinkKafkaProducer.SAFE_SCALE_DOWN_FACTOR.

        • Les transactions bloquent les lectures

          En mode read_committed, toute transaction qui n'a pas été validée ou abandonnée bloque les opérations de lecture sur l'ensemble du topic.

          Par exemple :

          • La transaction 1 écrit des données.

          • La transaction 2 écrit davantage de données et est validée.

          • Tant que la transaction 1 reste ouverte, les données de la transaction 2 validée sont invisibles pour les consommateurs.

          Cela a les implications suivantes :

          • Pendant le fonctionnement normal, la latence de visibilité des données est approximativement égale à l'intervalle de point de contrôle.

          • Si un job échoue, tout topic sur lequel il écrivait est bloqué pour les consommateurs jusqu'à ce que le job redémarre ou que la transaction expire. Dans les cas extrêmes, le processus d'expiration de la transaction lui-même peut également affecter les opérations de lecture.

        FAQ