- Publier ou s'abonner à des flux de données.
- Mettre en place un stockage tolérant aux pannes.
- Traiter les flux dès qu'ils sont disponibles.
Créer une table
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [ALIAS expr1],
name2 [type2] [ALIAS expr2],
...
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'host:port',
kafka_topic_list = 'topic1,topic2,...',
kafka_group_name = 'group_name',
kafka_format = 'data_format'[,]
[kafka_security_protocol = '',]
[kafka_sasl_mechanism = '',]
[kafka_sasl_username = '',]
[kafka_sasl_password = '',]
[kafka_autodetect_client_rack = '',]
[kafka_schema = '',]
[kafka_num_consumers = N,]
[kafka_max_block_size = 0,]
[kafka_skip_broken_messages = N,]
[kafka_commit_every_batch = 0,]
[kafka_client_id = '',]
[kafka_poll_timeout_ms = 0,]
[kafka_poll_max_batch_size = 0,]
[kafka_flush_interval_ms = 0,]
[kafka_consumer_reschedule_ms = 0,]
[kafka_thread_per_consumer = 0,]
[kafka_handle_error_mode = 'default',]
[kafka_commit_on_select = false,]
[kafka_consumer_acquire_timeout_ms = 30000,]
[kafka_max_rows_per_message = 1,]
[kafka_compression_codec = '',]
[kafka_compression_level = -1,]
[kafka_partition_shard_num = '',]
[kafka_shard_count = 0];Paramètres requis :
kafka_broker_list— Une liste de brokers séparés par des virgules (par exemple,localhost:9092).kafka_topic_list— Une liste de topics Kafka.kafka_group_name— Un groupe de consommateurs Kafka. Les offsets de lecture sont suivis séparément pour chaque groupe. Si vous ne souhaitez pas que les messages soient dupliqués dans le cluster, utilisez partout le même nom de groupe.kafka_format— Format des messages. Utilise la même notation que la fonction SQLFORMAT, par exempleJSONEachRow. Pour plus d'informations, consultez la section Formats.
Paramètres facultatifs :
kafka_security_protocol- Protocole utilisé pour communiquer avec les brokers. Valeurs possibles :plaintext,ssl,sasl_plaintext,sasl_ssl.kafka_sasl_mechanism- Mécanisme SASL à utiliser pour l’authentification. Valeurs possibles :GSSAPI,PLAIN,SCRAM-SHA-256,SCRAM-SHA-512,OAUTHBEARER,AWS_MSK_IAM.kafka_aws_region- Région AWS pour l’authentification MSK IAM. Détectée automatiquement à partir de l’adresse du broker si elle n’est pas spécifiée. Spécifiez-la explicitement lors de l’utilisation d’alias PrivateLink ou de noms d’hôte DNS personnalisés qui ne contiennent pas d’informations de région. Par défaut : vide (détection automatique).kafka_sasl_username- Nom d’utilisateur SASL à utiliser avec les mécanismesPLAINetSASL-SCRAM-...kafka_sasl_password- Mot de passe SASL à utiliser avec les mécanismesPLAINetSASL-SCRAM-...kafka_schema— Paramètre à utiliser si le format nécessite une définition de schéma. Par exemple, Cap'n Proto requiert le chemin vers le fichier de schéma ainsi que le nom de l’objet racineschema.capnp:Message.kafka_schema_registry_skip_bytes— Nombre d’octets à ignorer au début de chaque message lors de l’utilisation du registre de schémas avec des en-têtes d’enveloppe (par ex., AWS Glue Schema Registry, qui inclut une enveloppe de 19 octets). Plage :[0, 255]. Valeur par défaut :0.kafka_num_consumers— Nombre de consumers par table. Indiquez-en davantage si le throughput d’un consumer est insuffisant. Le nombre total de consumers ne doit pas dépasser le nombre de partitions du topic, puisqu’un seul consumer peut être attribué par partition, et ne doit pas non plus être supérieur au nombre de cœurs physiques du server sur lequel ClickHouse est déployé. Valeur par défaut :1.kafka_max_block_size— Taille maximale du batch (en messages) pour le poll. Valeur par défaut : max_insert_block_size.kafka_skip_broken_messages— Tolérance de l’analyseur de messages Kafka aux messages incompatibles avec le schéma par block. Sikafka_skip_broken_messages = N, le moteur ignore N messages Kafka qui ne peuvent pas être analysés (un message équivaut à une ligne de données). Valeur par défaut :0.kafka_commit_every_batch— Effectue un commit pour chaque batch consommé et traité, au lieu d’un seul commit après l’écriture d’un block entier. Valeur par défaut :0.kafka_client_id— Identifiant du client. Vide par défaut.kafka_poll_timeout_ms— Timeout pour un poll Kafka unique. Valeur par défaut : stream_poll_timeout_ms.kafka_poll_max_batch_size— Nombre maximal de messages pouvant être récupérés en un seul poll Kafka. Valeur par défaut : max_block_size.kafka_flush_interval_ms— Timeout pour le flush des données depuis Kafka. Valeur par défaut : stream_flush_interval_ms.kafka_consumer_reschedule_ms— Intervalle de replanification lorsque le stream processing Kafka est bloqué (par ex., lorsqu’aucun message n’est disponible à la consommation). Ce paramètre contrôle le délai avant que le consumer ne réessaie de poll. Ne doit pas dépasserkafka_consumers_pool_ttl_ms. Valeur par défaut :500millisecondes.kafka_thread_per_consumer— Fournit un thread indépendant pour chaque consumer. Lorsqu’elle est activée, chaque consumer flush les données indépendamment, en parallèle (sinon, les rows de plusieurs consumers sont squashed pour former un block). Valeur par défaut :0.kafka_handle_error_mode— Gestion des errors pour le Kafka engine. Valeurs possibles : default (une exception est levée si l’analyse d’un message échoue), stream (le message d’exception et le message brut sont enregistrés dans les virtual columns_erroret_raw_message), dead_letter_queue (les données liées à l’error sont enregistrées dans system.dead_letter_queue).kafka_commit_on_select— Effectue un commit des messages lorsqu’une requêteSELECTest exécutée. Valeur par défaut :false.kafka_consumer_acquire_timeout_ms— Timeout en millisecondes pour l’acquisition d’un consumer Kafka lors de requêtesSELECTdirectes sur une tableKafka2(avec stockage des offsets basé sur Keeper). Lorsque plusieurs requêtesSELECTdirectes concurrentes s’exécutent sur la même table, chacune doit attendre que des consumers deviennent disponibles. Ce timeout évite les deadlocks lorsque des requêtes détiennent différents sous-ensembles de consumers. Valeur par défaut :30000.kafka_max_rows_per_message— Nombre maximal de rows écrites dans un message Kafka pour les formats basés sur les lignes. Valeur par défaut :1.kafka_autodetect_client_rack— Définit automatiquement le paramètreclient.rackpourlibrdkafkaafin de privilégier les répliques Kafka les plus proches. Sources prises en charge :AWS_ZONE_IDpour l’ID de zone de disponibilité AWS IMDSv2, par exempleeuc1-az1;AWS_ZONE_NAMEpour le nom de zone de disponibilité AWS IMDSv2, par exempleeu-central-1a;GCP_ZONEpour la zone du service de métadonnées GCP, par exempleeurope-central2-a;CLICKHOUSEpour utiliser la détection interne de ClickHouse, qui peut s’appuyer sur les métadonnées du cloud ou sur la configuration ;AWS_ZONE_NAME_THEN_GCP_ZONEpour essayerAWS_ZONE_NAME, puisGCP_ZONE. Par défaut : chaîne vide, désactivé. Conseil : les formats de zone de disponibilité varient selon les environnements. Amazon MSK utilise généralement des ID de zone, privilégiez doncAWS_ZONE_ID. Confluent Cloud utilise généralement des noms de zone, privilégiez doncAWS_ZONE_NAME. En cas de doute, utilisezAWS_ZONE_NAME_THEN_GCP_ZONEou vérifiez la valeurbroker.racksur votre cluster. Remarque : les brokers Kafka doivent être configurés avecbroker.racketreplica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector.kafka_compression_codec— Codec de compression utilisé pour produire les messages. Valeurs prises en charge : chaîne vide,none,gzip,snappy,lz4,zstd. Si la chaîne est vide, le codec de compression n’est pas défini par la table ; les valeurs des fichiers de configuration ou la valeur par défaut delibrdkafkaseront alors utilisées. Par défaut : chaîne vide.kafka_compression_level— Paramètre de niveau de compression pour l’algorithme sélectionné par kafka_compression_codec. Des valeurs plus élevées offrent une meilleure compression, au prix d’une utilisation du CPU plus importante. La plage de valeurs utilisables dépend de l’algorithme :[0-9]pourgzip;[0-12]pourlz4; uniquement0poursnappy;[0-12]pourzstd;-1= niveau de compression par défaut dépendant du codec. Par défaut :-1.kafka_map_virtual_columns_on_write— Si cette option est activée, les colonnes portant les noms spéciaux_key,_timestamp,_headers.nameet_headers.valuedans le schéma de la table sont associées aux métadonnées correspondantes du message Kafka lors deINSERTet sont exclues du corps du message. Voir Correspondance entre les colonnes et les métadonnées des messages Kafka. Par défaut :false.kafka_partition_shard_num— Numéro du shard actuel pour l’affinité statique entre les partitions et les shards. Doit être compris entre 1 etkafka_shard_count, bornes incluses. Les partitions sont attribuées selon la formulepartition_id % kafka_shard_count == kafka_partition_shard_num - 1. Prend en charge l’expansion de macros (par ex.,'{shard}'). Doit être utilisé aveckafka_shard_count. Pris en charge uniquement avec StorageKafka2 (nécessitekafka_keeper_pathetkafka_replica_name). Valeur par défaut :''(désactivé).kafka_shard_count— Nombre total de shards participant à la consommation. Utilisé aveckafka_partition_shard_numpour attribuer statiquement les partitions. Doit être utilisé aveckafka_partition_shard_num. Pris en charge uniquement avec StorageKafka2. Valeur par défaut :0(désactivé).
Exemples :
CREATE TABLE queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka('localhost:9092', 'topic', 'group1', 'JSONEachRow');
SELECT * FROM queue LIMIT 5;
CREATE TABLE queue2 (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka SETTINGS kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
CREATE TABLE queue3 (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka('localhost:9092', 'topic', 'group1')
SETTINGS kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;Méthode obsolète pour créer une table
Kafka(kafka_broker_list, kafka_topic_list, kafka_group_name, kafka_format
[, kafka_row_delimiter, kafka_schema, kafka_num_consumers, kafka_max_block_size, kafka_skip_broken_messages, kafka_commit_every_batch, kafka_client_id, kafka_poll_timeout_ms, kafka_poll_max_batch_size, kafka_flush_interval_ms, kafka_consumer_reschedule_ms, kafka_thread_per_consumer, kafka_handle_error_mode, kafka_commit_on_select, kafka_max_rows_per_message]);Description
Les messages livrés sont suivis automatiquement, de sorte que chaque message d’un groupe n’est comptabilisé qu’une seule fois. Si vous souhaitez obtenir les données deux fois, créez une copie de la table avec un autre nom de groupe.
Les groupes sont flexibles et synchronisés sur le cluster. Par exemple, si vous avez 10 topics et 5 copies d’une table dans un cluster, chaque copie reçoit 2 topics. Si le nombre de copies change, les topics sont automatiquement redistribués entre les copies. Pour en savoir plus à ce sujet, consultez http://kafka.apache.org/intro.
Il est recommandé que chaque topic Kafka dispose de son propre groupe de consommateurs dédié, afin de garantir une association exclusive entre le topic et le groupe, en particulier dans les environnements où les topics peuvent être créés et supprimés dynamiquement (par exemple, en test ou en préproduction).
SELECT n’est pas particulièrement utile pour lire les messages (sauf pour le débogage), car chaque message ne peut être lu qu’une seule fois. Il est plus pratique de créer des flux en temps réel à l’aide de vues matérialisées. Pour cela :
- Utilisez le moteur pour créer un consommateur Kafka et considérez-le comme un flux de données.
- Créez une table avec la structure souhaitée.
- Créez une vue matérialisée qui convertit les données du moteur et les insère dans une table créée précédemment.
Lorsque la MATERIALIZED VIEW se connecte au moteur, elle commence à collecter les données en arrière-plan. Cela vous permet de recevoir en continu des messages de Kafka et de les convertir au format requis à l’aide de SELECT.
Une table Kafka peut avoir autant de vues matérialisées que vous le souhaitez : elles ne lisent pas directement les données de la table Kafka, mais reçoivent les nouveaux enregistrements (par blocs). Vous pouvez ainsi écrire dans plusieurs tables avec des niveaux de granularité différents (avec regroupement - agrégation ou sans).
Exemple, en utilisant des collections nommées pour stocker les paramètres de connexion :
CREATE NAMED COLLECTION kafka_creds AS
kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow';
CREATE TABLE queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka(kafka_creds);
CREATE TABLE daily (
day Date,
level String,
total UInt64
) ENGINE = SummingMergeTree
PARTITION BY toYYYYMM(day)
ORDER BY (day, level);
CREATE MATERIALIZED VIEW consumer TO daily
AS SELECT toDate(toDateTime(timestamp)) AS day, level, count() AS total
FROM queue GROUP BY day, level;
SELECT level, sum(total) FROM daily GROUP BY level;Pour améliorer les performances, les messages reçus sont regroupés en blocs de la taille de max_insert_block_size. Si le bloc n'a pas été constitué dans le délai de stream_flush_interval_ms millisecondes, les données seront écrites dans la table, même si le bloc n'est pas complet.
Pour arrêter de recevoir les données du topic ou modifier la logique de conversion, détachez la vue matérialisée :
DETACH TABLE consumer;
ATTACH TABLE consumer;Si vous souhaitez modifier la table cible en utilisant ALTER, nous vous recommandons de désactiver la vue matérialisée afin d’éviter tout écart entre la table cible et les données de la vue.
Configuration
Comme pour GraphiteMergeTree, le moteur Kafka prend en charge une configuration étendue via le fichier de configuration de ClickHouse. Vous pouvez utiliser deux clés de configuration : une globale (sous <kafka>) et une au niveau du topic (sous <kafka><kafka_topic>). La configuration globale est appliquée en premier, puis la configuration au niveau du topic est appliquée (si elle existe).
<kafka>
<!-- Global configuration options for all tables of Kafka engine type -->
<debug>cgrp</debug>
<statistics_interval_ms>3000</statistics_interval_ms>
<kafka_topic>
<name>logs</name>
<statistics_interval_ms>4000</statistics_interval_ms>
</kafka_topic>
<!-- Settings for consumer -->
<consumer>
<auto_offset_reset>smallest</auto_offset_reset>
<kafka_topic>
<name>logs</name>
<fetch_min_bytes>100000</fetch_min_bytes>
</kafka_topic>
<kafka_topic>
<name>stats</name>
<fetch_min_bytes>50000</fetch_min_bytes>
</kafka_topic>
</consumer>
<!-- Settings for producer -->
<producer>
<kafka_topic>
<name>logs</name>
<retry_backoff_ms>250</retry_backoff_ms>
</kafka_topic>
<kafka_topic>
<name>stats</name>
<retry_backoff_ms>400</retry_backoff_ms>
</kafka_topic>
</producer>
</kafka>Pour obtenir la liste des options de configuration possibles, consultez la référence de configuration de librdkafka. Dans la configuration de ClickHouse, utilisez le caractère de soulignement (_) à la place d’un point. Par exemple, check.crcs=true deviendra <check_crcs>true</check_crcs>.
Authentification IAM pour AWS MSK
AWS MSK prend en charge l'authentification basée sur IAM, ce qui permet de se connecter à des clusters Kafka à l'aide d'identifiants AWS au lieu de gérer des noms d'utilisateur et mots de passe distincts.
Configuration de base :
Définissez kafka_sasl_mechanism = 'AWS_MSK_IAM' dans les paramètres de votre table :
CREATE TABLE msk_queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'b-1.mycluster.kafka.us-east-1.amazonaws.com:9098',
kafka_topic_list = 'my-topic',
kafka_group_name = 'my-group',
kafka_format = 'JSONEachRow',
kafka_sasl_mechanism = 'AWS_MSK_IAM';La région AWS est automatiquement extraite de l’endpoint du broker par correspondance de motifs :
- MSK provisionné :
b-X.cluster.kafka.<region>.amazonaws.com:9098 - MSK serverless :
boot-X.kafka-serverless.<region>.amazonaws.com:9098 - VPC Endpoint :
vpce-X.kafka.<region>.vpce.amazonaws.com:9098
Identifiants AWS :
Les identifiants sont toujours chargés depuis ~/.aws/credentials et ~/.aws/config (fichiers de profil AWS) lorsqu’ils sont présents. Pour activer aussi les profils d’instance EC2, les variables d’environnement (AWS_ACCESS_KEY_ID, etc.), les rôles de tâche ECS et les autres sources automatiques d’identifiants, ajoutez ce qui suit à votre configuration du serveur :
<kafka>
<use_environment_credentials>true</use_environment_credentials>
</kafka>Ce paramètre ne peut être configuré que par les administrateurs du serveur. Valeur par défaut : false.
PrivateLink et DNS personnalisé :
Lorsque vous utilisez des alias PrivateLink ou des noms d’hôte DNS personnalisés qui ne contiennent pas d’informations sur la région, indiquez explicitement la région AWS :
CREATE TABLE msk_privatelink_queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'my-privatelink-alias.internal.example.com:9098',
kafka_topic_list = 'my-topic',
kafka_group_name = 'my-group',
kafka_format = 'JSONEachRow',
kafka_sasl_mechanism = 'AWS_MSK_IAM',
kafka_aws_region = 'us-east-1';Autorisations IAM :
Autorisations du consommateur (pour lire les messages) :
{
"Version": "2012-10-17",
"Statement": [{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect",
"kafka-cluster:DescribeTopic",
"kafka-cluster:ReadData",
"kafka-cluster:AlterGroup",
"kafka-cluster:DescribeGroup"
],
"Resource": [
"arn:aws:kafka:REGION:ACCOUNT:cluster/CLUSTER_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:topic/CLUSTER_NAME/TOPIC_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:group/CLUSTER_NAME/CONSUMER_GROUP/*"
]
}]
}Autorisations du producteur (pour écrire des messages) :
{
"Version": "2012-10-17",
"Statement": [{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect",
"kafka-cluster:DescribeTopic",
"kafka-cluster:WriteData"
],
"Resource": [
"arn:aws:kafka:REGION:ACCOUNT:cluster/CLUSTER_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:topic/CLUSTER_NAME/TOPIC_NAME/*"
]
}]
}Prise en charge de Kerberos
Pour utiliser Kafka avec Kerberos, ajoutez l’élément enfant security_protocol avec la valeur sasl_plaintext. Cela suffit si le ticket d’octroi de tickets Kerberos est obtenu et mis en cache par les mécanismes du système d’exploitation.
ClickHouse peut gérer les informations d’identification Kerberos à l’aide d’un fichier keytab. Tenez compte des éléments enfants sasl_kerberos_service_name, sasl_kerberos_keytab et sasl_kerberos_principal.
Exemple :
<!-- Kerberos-aware Kafka -->
<kafka>
<security_protocol>SASL_PLAINTEXT</security_protocol>
<sasl_kerberos_keytab>/home/kafkauser/kafkauser.keytab</sasl_kerberos_keytab>
<sasl_kerberos_principal>kafkauser/kafkahost@EXAMPLE.COM</sasl_kerberos_principal>
</kafka>Colonnes virtuelles
_topic— Topic Kafka. Type de données :LowCardinality(String)._key— Clé du message. Type de données :String._offset— Offset du message. Type de données :UInt64._timestamp— Horodatage du message. Type de données :Nullable(DateTime)._timestamp_ms— Horodatage du message en millisecondes. Type de données :Nullable(DateTime64(3))._partition— Partition du topic Kafka. Type de données :UInt64._headers.name— Tableau des clés d'en-tête du message. Type de données :Array(String)._headers.value— Tableau des valeurs d'en-tête du message. Type de données :Array(String).
Colonnes virtuelles supplémentaires lorsque kafka_handle_error_mode='stream' :
_raw_message- Message brut qui n'a pas pu être analysé correctement. Type de données :String._error- Message d'exception généré lors d'une erreur d'analyse. Type de données :String.
Remarque : les colonnes virtuelles _raw_message et _error sont renseignées uniquement en cas d'exception pendant l'analyse ; elles sont toujours vides lorsque le message a été analysé correctement.
Correspondance entre les colonnes et les métadonnées des messages Kafka
Lors de la production de messages avec INSERT INTO, le moteur Kafka utilise toujours une colonne nommée _key (de type String) comme clé du message Kafka et une colonne nommée _timestamp (de type DateTime) comme horodatage du message Kafka — si ces colonnes existent dans la table. Par défaut, ces colonnes apparaissent également dans la charge utile du message produit, aux côtés des autres colonnes.
Avec kafka_map_virtual_columns_on_write = 1, le comportement change :
_key(typeString) — associé à la clé du message Kafka._timestamp(typeDateTime) — associé à l’horodatage du message Kafka._headers.name(typeArray(String)) et_headers.value(typeArray(String)) — associés aux en-têtes des messages Kafka. Chaque paire(_headers.name[i], _headers.value[i])devient un en-tête Kafka. Comme_headers.nameet_headers.valuepartagent le préfixe Nested_headers, ClickHouse exige que les deux tableaux aient la même taille pour chaque ligne.
Les colonnes portant ces noms sont exclues de la charge utile du message uniquement si leurs types correspondent à ceux indiqués ci-dessus ; sinon, elles restent dans la charge utile, de sorte que les schémas qui réutilisent ces noms par hasard pour des données sans rapport continuent de fonctionner.
Exemple :
CREATE TABLE kafka_out
(
event_json String,
`_key` String,
`_timestamp` DateTime,
`_headers.name` Array(String),
`_headers.value` Array(String)
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'broker:9092',
kafka_topic_list = 'events',
kafka_group_name = 'events-producer',
kafka_format = 'JSONEachRow',
kafka_map_virtual_columns_on_write = 1;
INSERT INTO kafka_out VALUES
('{\"a\":1}', 'session-42', now(), ['source', 'trace_id'], ['api', 'abc-123']);Le message Kafka produit contient la charge utile {"event_json":"{\"a\":1}"}, la clé session-42, l’horodatage actuel et deux en-têtes source=api et trace_id=abc-123.
Prise en charge des formats de données
Le moteur Kafka prend en charge tous les formats pris en charge par ClickHouse. Le nombre de lignes dans un message Kafka dépend du fait que le format soit basé sur les lignes ou sur les blocs :
- Pour les formats basés sur les lignes, le nombre de lignes dans un message Kafka peut être contrôlé en définissant
kafka_max_rows_per_message. - Pour les formats basés sur les blocs, il n’est pas possible de diviser un bloc en parties plus petites, mais le nombre de lignes dans un bloc peut être contrôlé par le paramètre général max_block_size.
Moteur pour stocker les offsets validés dans ClickHouse Keeper
Si allow_experimental_kafka_offsets_storage_in_keeper est activé, deux paramètres supplémentaires peuvent être spécifiés pour le moteur de table Kafka :
kafka_keeper_pathspécifie le chemin vers la table dans ClickHouse Keeperkafka_replica_namespécifie le nom de la réplique dans ClickHouse Keeper
Soit les deux paramètres doivent être spécifiés, soit aucun des deux. Lorsque les deux sont spécifiés, un nouveau moteur Kafka expérimental est utilisé. Ce nouveau moteur ne dépend pas du stockage des offsets validés dans Kafka, mais les stocke dans ClickHouse Keeper. Il essaie toujours de valider les offsets dans Kafka, mais il ne dépend de ces offsets qu’au moment de la création de la table. Dans tous les autres cas (si la table est redémarrée ou restaurée après une erreur), les offsets stockés dans ClickHouse Keeper sont utilisés pour reprendre la consommation des messages. En plus de l’offset validé, il stocke également le nombre de messages consommés dans le dernier lot, de sorte que, si l’insert échoue, le même nombre de messages sera consommé, ce qui permet la déduplication si nécessaire.
Affinité statique entre partitions et segments
Lorsque vous utilisez StorageKafka2, vous pouvez activer une affinité statique entre partitions et segments en spécifiant kafka_partition_shard_num et kafka_shard_count. Cela permet à plusieurs instances ClickHouse (segments) de consommer le même topic Kafka, chaque segment ne traitant qu’un sous-ensemble déterministe de partitions selon la formule :
partition_id % kafka_shard_count == kafka_partition_shard_num - 1Les deux paramètres doivent être spécifiés ensemble ; n’en spécifier un seul déclenche une exception. La valeur de kafka_partition_shard_num doit être comprise entre 1 et kafka_shard_count, bornes incluses. Elle prend en charge l’expansion de macros (par exemple, '{shard}'), réévaluée à chaque démarrage du serveur. Cela permet de partager les mêmes métadonnées de table entre les segments d’une base de données Replicated, chaque segment résolvant sa propre valeur. La validation est effectuée après l’expansion des macros.
Tous les segments doivent utiliser la même valeur pour kafka_keeper_path. Toutes les répliques partagent les offsets validés et les tailles d’intention, mais seules les répliques ayant le même numéro de segment sont en concurrence pour les verrous de partition (à condition que toutes les répliques aient le même nombre de segments).
Exemple avec 3 segments consommant un topic comportant 12 partitions :
-- Shard 1: consumes partitions 0, 3, 6, 9
CREATE TABLE kafka_shard1 (key UInt64, value String)
ENGINE = Kafka('localhost:9092', 'my-topic', 'my-group', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/kafka/{database}',
kafka_replica_name = '{replica}',
kafka_partition_shard_num = '1',
kafka_shard_count = 3
SETTINGS allow_experimental_kafka_offsets_storage_in_keeper = 1;
-- Shard 2: consumes partitions 1, 4, 7, 10
CREATE TABLE kafka_shard2 (key UInt64, value String)
ENGINE = Kafka('localhost:9092', 'my-topic', 'my-group', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/kafka/{database}',
kafka_replica_name = '{replica}',
kafka_partition_shard_num = '2',
kafka_shard_count = 3
SETTINGS allow_experimental_kafka_offsets_storage_in_keeper = 1;Combiné à des répliques (plusieurs valeurs de kafka_replica_name partageant le même kafka_keeper_path), le filtre d’affinité est d’abord appliqué pour déterminer les partitions éligibles, puis les verrous ZooKeeper répartissent ces dernières entre les répliques.
Exemple :
CREATE TABLE experimental_kafka (key UInt64, value UInt64)
ENGINE = Kafka('localhost:19092', 'my-topic', 'my-consumer', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/{database}/{uuid}',
kafka_replica_name = '{replica}'
SETTINGS allow_experimental_kafka_offsets_storage_in_keeper=1;Limites connues
Comme le nouveau moteur est expérimental, il n’est pas encore prêt pour une utilisation en production. L’implémentation présente quelques limites connues :
- Supprimer puis recréer rapidement la table, ou spécifier le même chemin ClickHouse Keeper pour différents moteurs, peut entraîner des problèmes. Comme bonne pratique, vous pouvez utiliser
{uuid}danskafka_keeper_pathpour éviter les conflits de chemins. - Pour garantir des lectures répétables, les messages ne peuvent pas être consommés à partir de plusieurs partitions sur un seul thread. En revanche, les consommateurs Kafka doivent être interrogés régulièrement pour être maintenus actifs. En raison de ces deux contraintes, nous avons décidé d’autoriser la création de plusieurs consommateurs uniquement si
kafka_thread_per_consumerest activé, sinon il est trop compliqué d’éviter les problèmes liés à l’interrogation régulière des consommateurs. - Lors de l’utilisation de l’affinité de partition, tous les segments doivent utiliser le même
kafka_shard_count; sinon, certaines partitions peuvent être consommées par plusieurs segments ou rester non consommées.
Voir aussi