Général
|
Option |
Description |
Type |
Obligatoire |
Valeur par défaut |
Remarques |
|
connector |
Type de connecteur. |
String |
Oui |
– |
La valeur doit être |
|
properties.bootstrap.servers |
Liste des adresses des courtiers Kafka. |
String |
Oui |
– |
Format : |
|
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 Le connecteur Kafka écrase ces options ; vous ne pouvez donc pas les configurer de cette manière :
|
|
format |
Format de sérialisation et de désérialisation de la valeur d'un message Kafka. |
String |
Non |
– |
Formats pris en charge :
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 :
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, |
|
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.format |
Format de sérialisation et de désérialisation de la valeur d'un message Kafka. |
String |
Non |
– |
Cette configuration équivaut à |
|
value.fields-include |
Définit si les champs de clé sont inclus dans le format de valeur. |
String |
Non |
ALL |
Valeurs valides :
|
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 Remarque
Vous pouvez spécifier cette option ou |
|
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 :
Remarque
Vous pouvez spécifier cette option ou |
|
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 |
|
scan.startup.mode |
Décalage de démarrage du consommateur Kafka. |
String |
Non |
group-offsets |
Valeurs valides :
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 |
String |
Non |
– |
Par exemple, |
|
scan.startup.timestamp-millis |
Horodatage de démarrage en millisecondes lorsque |
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, Remarque
|
|
scan.check.duplicated.group.id |
Vérifie si un autre consommateur actif utilise déjà le |
Boolean |
Non |
false |
Valeurs valides :
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 :
|
|
sink.delivery-guarantee |
Garantie de livraison du sink. |
String |
Non |
at-least-once |
Valeurs valides :
Remarque
Lors de l'utilisation de la sémantique |
|
sink.transactional-id-prefix |
Préfixe d'ID de transaction. Obligatoire lorsque |
String |
Oui, si |
– |
Requis uniquement lorsque sink.delivery-guarantee est défini sur |
|
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 |
|
|
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.
ImportantLimitations 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.idempotenceest définie par défaut surtrue. 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 :
Le client utilise les adresses spécifiées dans
bootstrap.serverspour établir une connexion initiale au cluster Kafka.Le cluster Kafka renvoie les métadonnées de chaque courtier, y compris leurs endpoints.
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
-
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. -
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.
-
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
-
Utilisez la fonctionnalité Network Probe
Cette fonctionnalité vous aide à écarter les problèmes de connectivité avec l'adresse
bootstrap.serverset à vérifier que l'endpoint interne ou public correct est utilisé. -
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.
-
Vérifiez la configuration
Utilisez l'outil zkCli.sh ou zookeeper-shell.sh pour vous connecter au cluster ZooKeeper utilisé par Kafka.
-
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} -
Utilisez la fonctionnalité Network Probe dans la console de développement Realtime Compute for Apache Flink pour tester l'accessibilité de cette adresse.
RemarqueSi l'adresse n'est pas accessible, contactez vos administrateurs Kafka pour vérifier et corriger les configurations
listenersetadvertised.listenersafin 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é.
-
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 :
|
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.locationetproperties.ssl.keystore.locationavec 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.
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.
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;
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
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 estmy-groupet 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).KafkaRecordDeserializationSchemadéfinit comment analyser unConsumerRecordKafka. 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.DeserializationSchemadé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));
RemarquePour analyser un
ConsumerRecordcomplet, vous devez implémenter l'interfaceKafkaRecordDeserializationSchema.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())RemarqueSi 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.RemarqueUne 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 :setUnboundedpour le mode streaming etsetBoundedpour 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.ImportantLa 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.
RemarqueSi 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
jsonest pris en charge. -
Pour le puits, les valeurs valides sont :
-
csv
-
json
-
RemarqueCette 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-jsonetcanal-jsonnécessitent Realtime Compute for Apache Flink (VVR) version 8.0.10 ou ultérieure. -
Le format
jsonné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.RemarqueSpé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 paruser_event_. -
prod\.logs\..*: Correspond aux topics préfixés parprod.logs.(le caractère.doit être échappé).
RemarqueSpé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
earliestoulatestafin 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.
RemarqueCe 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.modeest défini surspecific-offsets.Non
String
–
Par exemple,
partition:0,offset:42;partition:1,offset:300scan.startup.timestamp-millis
Horodatage de démarrage en millisecondes lorsque
scan.startup.modeest défini surtimestamp.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.idest 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-
Pour plus d'informations sur l'analyse des schémas, consultez Politiques d'analyse et d'évolution des schémas.
-
Cette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 8.0.11.
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é devientkey_a.RemarqueLa valeur de
key.fields-prefixne peut pas être un préfixe de la valeur devalue.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é devientvalue_b.RemarqueLa valeur de
value.fields-prefixne peut pas être un préfixe de la valeur dekey.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,headersetleader-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, utilisezCREATE 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.
RemarqueCette 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
RemarqueCette 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-errorsest défini surtrue.Non
Integer
-1
Ce paramètre s'applique uniquement lorsque
ingestion.ignore-errorsest défini surtrue. 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.RemarqueCette 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.
RemarqueCette 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
truesi les données d'une seule table Debezium JSON sont distribuées sur plusieurs partitions.RemarqueCette option de configuration n'est prise en charge qu'à partir de Ververica Runtime (VVR) 8.0.11.
ImportantLa 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.enabledans 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
Stringlors 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.RemarqueCe 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.
RemarqueCette option de configuration est prise en charge uniquement dans Ververica Runtime (VVR) 8.0.11 et versions ultérieures.
ImportantLa 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
databasepré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
tablepré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
Stringlors 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.RemarqueCe 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
sqlTypepré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 champsqlType. -
MYSQL_TYPE : analyse le schéma à partir du tableau
mysqlTypepré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_TYPEest 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
TIMESTAMPau type CDCTIMESTAMP.-
true : le type MySQL
TIMESTAMPest mappé au type CDCTIMESTAMP. -
false : le type MySQL
TIMESTAMPest mappé au type CDCTIMESTAMP_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 MySQLTINYINT(1)doit être mappé au type CDCBOOLEAN.-
true : le type MySQL
TINYINT(1)est mappé au type CDCBOOLEAN. -
false : le type MySQL
TINYINT(1)est mappé au type CDCTINYINT(1).
Cette option s'applique uniquement lorsque
canal-json.infer-schema.strategyest défini surMYSQL_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 exemple2020-12-30 12:13:14.123. -
ISO-8601 : analyse les horodatages d'entrée au format
yyyy-MM-ddTHH:mm:ss.s{precision}, par exemple2020-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
Stringlors 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.RemarqueCe 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 champidest de type BIGINT et que le champnameest 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.
RemarqueCette 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.
RemarqueSi 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 exempledatabaseName.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.RemarqueCe 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
oldd'un message UPDATE au format Canal JSON contient uniquement les anciennes valeurs des champs qui ont changé.RemarqueCe 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
dbSchema
schemaTable
tableSi ce paramètre est désactivé, le mappage est le suivant :
CDC Table ID Part
Debezium JSON Key
Namespace
Non mappé
Schema
dbTable
tableRemarqueCe 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
routespécifie le topic Kafka de destination pour la table source.
RemarquePar 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.commitetauto.commit.interval.ms.RemarqueLa 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)etsetProperty(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
-1désactive la découverte dynamique des partitions.RemarqueEn 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.
ImportantPour garantir un fonctionnement correct, le connecteur DataStream Connector de Kafka écrase les propriétés configurées manuellement suivantes :
-
key.deserializerest toujours remplacé par org.apache.kafka.common.serialization.ByteArrayDeserializer. -
value.deserializerest toujours remplacé par org.apache.kafka.common.serialization.ByteArrayDeserializer. -
auto.offset.reset.strategyest remplacé par la stratégie fournie parOffsetsInitializer.
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-groupeKafkaSourceReader.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-totalest enregistrée à l'emplacement .operator.KafkaSourceReader.KafkaConsumer.records-consumed-total.Utilisez la propriété
register.consumer.metricspour 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.serversest requise. Elle spécifie une liste de brokers Kafka séparés par des virgules.Sérialiseur d'enregistrements
Vous devez fournir un
KafkaRecordSerializationSchemapour convertir les données d'entrée en unProducerRecordKafka. 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
ProducerRecordoffre 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.serversest 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.
RemarqueLors 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:9092Straté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 :
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.
RemarquePour 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.
ImportantLa 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.
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
idcommeBIGINTet de la colonnenamecommeVARCHAR(10)pour la tabledb1.t1, et le type initial de la colonneidcommeBIGINTpour la tabledb1.t2.Les instructions DDL utilisent la syntaxe Flink SQL.
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: 0Cette configuration spécifie que tous les champs
idsont de typeBIGINTet que tous les champsnamesont de typeVARCHAR(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,col2Avec cette configuration, le connecteur envoie chaque enregistrement à une table où le nom de la base de données est la valeur du champ
col1et le nom de la table est la valeur du champcol2.
-
Informations sur la clé primaire
Pour le format Canal JSON, le champ
pkNamesdans 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.

-
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
INTversBIGINT) 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 NULLversNULLABLE.
-
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 surSQL_TYPEpour utiliser les types du champsqlType. 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: 1000Avec 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: -1Bien 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: loggerMappage 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.mytable1etmydb.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.mytable1sont écrites dans un topic nommémydb.mytable1, et les données demydb.mytable2sont é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: mytableDans ce cas, toutes les données de
mydb.mytable1etmydb.mytable2sont é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.mytable1etmydb.mytable2sont écrites dans le topicmytable, et le nom de la table dans les messages Kafka (au format Debezium JSON ou Canal JSON) est conservé en tant quemydb.mytable1oumydb.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
customerset é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: truePour 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 :
-
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 -
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
customerset ajoute les colonnes de métadonnéestopicetpartition: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: trueGé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: loggerL'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: loggerLire 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
dbettblcomme 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: loggerSpé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: loggerLire 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: loggerVous 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: loggerLire des données Debezium JSON
L'exemple suivant lit le topic
customerset é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: trueSynchroniser 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-zeroSé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.mspour un broker Kafka est de 15 minutes.Par défaut, Flink Kafka Sink définit le paramètre
transaction.timeout.mssur 1 heure.Vous devez augmenter la valeur de
transaction.max.timeout.mssur 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_ONCEutilise 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
-