Kafka vers ClickHouse
Si vous utilisez ClickHouse Cloud, nous vous recommandons plutôt d’utiliser ClickPipes. ClickPipes prend nativement en charge les connexions sur réseau privé, la mise à l’échelle indépendante de l’ingestion et des ressources du cluster, ainsi qu’un monitoring complet pour l’ingestion continue de données Kafka dans ClickHouse.
Vue d’ensemble
Étapes
1
Préparez
Si vous avez des données dans un topic cible, vous pouvez adapter ce qui suit à votre jeu de données. Sinon, un jeu de données d’exemple GitHub est fourni ici. Ce jeu de données est utilisé dans les exemples ci-dessous et repose sur un schéma réduit ainsi que sur un sous-ensemble des lignes (plus précisément, nous nous limitons aux événements GitHub concernant le dépôt ClickHouse), par souci de concision, par rapport au jeu de données complet disponible ici. Cela reste néanmoins suffisant pour que la plupart des requêtes publiées avec le jeu de données fonctionnent.
2
Configurer ClickHouse
Cette étape est requise si vous vous connectez à un Kafka sécurisé. Ces paramètres ne peuvent pas être fournis via les commandes SQL DDL et doivent être configurés dans le Placez l’extrait ci-dessus soit dans un nouveau fichier de votre répertoire conf.d/, soit fusionnez-le dans des fichiers de configuration existants. Pour les paramètres qui peuvent être configurés, consultez ici.Nous allons également créer une base de données appelée Une fois la base de données créée, vous devrez ensuite la sélectionner :
config.xml de ClickHouse. Nous supposons que vous vous connectez à une instance sécurisée par SASL. C’est la méthode la plus simple pour interagir avec Confluent Cloud.KafkaEngine à utiliser dans ce tutoriel :3
Créez la table de destination
Préparez votre table de destination. Dans l’exemple ci-dessous, nous utilisons un schéma GitHub simplifié par souci de concision. Notez que, bien que nous utilisions un moteur de table MergeTree, cet exemple pourrait facilement être adapté à n’importe quel membre de la famille MergeTree.
4
Créez le topic et peuplez-le
Ensuite, nous allons créer un topic. Il existe plusieurs outils pour le faire. Si Kafka s’exécute localement sur notre machine ou dans un conteneur Docker, RPK fonctionne bien. Nous pouvons créer un topic appelé Si nous exécutons Kafka sur Confluent Cloud, nous préférerons peut-être utiliser la CLI Confluent :Nous devons maintenant alimenter ce topic avec des données, ce que nous ferons à l’aide de kcat. Nous pouvons exécuter une commande semblable à la suivante si Kafka s’exécute localement avec l’authentification désactivée :Ou bien ce qui suit si notre cluster Kafka utilise SASL pour l’authentification :Le jeu de données contient 200 000 lignes, il devrait donc être ingéré en quelques secondes seulement. Si vous souhaitez travailler avec un jeu de données plus volumineux, consultez la section sur les grands jeux de données du dépôt GitHub ClickHouse/kafka-samples.
github avec 5 partitions en exécutant la commande suivante :5
Créez le moteur de table Kafka
L’exemple ci-dessous crée un moteur de table avec le même schéma que la table MergeTree. Ce n’est pas strictement nécessaire, car vous pouvez avoir des colonnes alias ou éphémères dans la table cible. Les paramètres sont toutefois importants ; notez l’utilisation de Nous abordons ci-dessous les paramètres du moteur et l’optimisation des performances. À ce stade, une simple requête select sur la table
JSONEachRow comme type de données pour consommer du JSON depuis un topic Kafka. Les valeurs github et clickhouse représentent respectivement le nom du topic et du groupe de consommateurs. Les topics peuvent en réalité être une liste de valeurs.github_queue devrait lire quelques lignes. Notez que cela fera avancer les offsets du consumer, empêchant de relire ces lignes sans réinitialisation. Notez la limite et le paramètre requis stream_like_engine_allow_direct_select.6
Créez la vue matérialisée
La vue matérialisée reliera les deux tables créées précédemment, en lisant les données du moteur de table Kafka et en les insérant dans la table MergeTree cible. Nous pouvons effectuer un certain nombre de transformations des données. Nous ferons une simple lecture et insertion. L’utilisation de * suppose que les noms de colonnes sont identiques (respect de la casse).Au moment de sa création, la vue matérialisée se connecte au moteur Kafka et commence à lire, en insérant des lignes dans la table de destination. Ce processus se poursuivra indéfiniment, les messages insérés par la suite dans Kafka étant consommés. N’hésitez pas à réexécuter le script d’insertion pour insérer d’autres messages dans Kafka.
7
Vérifiez que des lignes ont bien été insérées
Vérifiez que des données sont présentes dans la table cible :Vous devriez voir 200 000 lignes :
Opérations courantes
Arrêt et redémarrage de la consommation des messages
Ajout des métadonnées Kafka
_.
Vous trouverez une liste complète des colonnes virtuelles ici.
Pour mettre à jour notre table avec les colonnes virtuelles, nous devrons supprimer la vue matérialisée, rattacher la table du moteur Kafka, puis recréer la vue matérialisée.
Modifier les paramètres du moteur Kafka
Débogage des problèmes
Gestion des messages mal formés
- Traitez le champ de message comme une chaîne de caractères. Des fonctions peuvent être utilisées dans l’instruction de vue matérialisée pour effectuer le nettoyage et le transtypage si nécessaire. Cela ne constitue pas une solution de production, mais peut aider pour une ingestion ponctuelle.
- Si vous consommez du JSON depuis un topic avec le format JSONEachRow, utilisez le paramètre
input_format_skip_unknown_fields. Lors de l’écriture des données, ClickHouse lève par défaut une exception si les données d’entrée contiennent des colonnes qui n’existent pas dans la table cible. En revanche, si cette option est activée, ces colonnes supplémentaires seront ignorées. Là encore, ce n’est pas une solution de niveau production et cela pourrait prêter à confusion. - Envisagez le paramètre
kafka_skip_broken_messages. Il oblige l’utilisateur à spécifier un niveau de tolérance par bloc pour les messages mal formés, en tenant compte de kafka_max_block_size. Si cette tolérance est dépassée (mesurée en nombre absolu de messages), le comportement habituel reprendra avec une exception, et les autres messages seront ignorés.
Sémantique de livraison et problèmes liés aux doublons
Insertions avec quorum
ClickHouse vers Kafka
Étapes
1
Insertion directe de lignes
Commencez par vérifier le nombre de lignes dans la table cible.Vous devriez avoir 200 000 lignes :Réinsérez maintenant des lignes de la table cible GitHub dans le moteur de table Kafka github_queue. Notez que nous utilisons le format JSONEachRow et que nous limitons le SELECT à 100.Recomptez les lignes de GitHub pour confirmer que leur nombre a augmenté de 100. Comme le montre le schéma ci-dessus, des lignes ont été insérées dans Kafka via le moteur de table Kafka avant d’être relues par ce même moteur et insérées dans la table cible GitHub par notre vue matérialisée !Vous devriez voir 100 lignes supplémentaires :
2
Utilisation de vues matérialisées
Nous pouvons utiliser des vues matérialisées pour envoyer des messages vers un moteur Kafka (et un topic) lorsque des documents sont insérés dans une table. Lorsque des lignes sont insérées dans la table GitHub, une vue matérialisée est déclenchée, ce qui entraîne l’insertion de ces lignes dans un moteur Kafka puis dans un nouveau topic. Là encore, le schéma suivant l’illustre le mieux :Créez un nouveau topic Kafka Créez maintenant une nouvelle vue matérialisée Si vous insérez des documents dans le topic github d’origine, créé dans le cadre de Kafka to ClickHouse, ils apparaîtront comme par magie dans le topic “github_clickhouse”. Vérifiez-le à l’aide des outils natifs de Kafka. Par exemple, ci-dessous, nous insérons 100 lignes dans le topic github à l’aide de kcat pour un topic hébergé sur Confluent Cloud :Une lecture sur le topic Bien qu’il s’agisse d’un exemple élaboré, celui-ci illustre toute la puissance des vues matérialisées utilisées conjointement avec le moteur Kafka.
github_out ou un équivalent. Assurez-vous qu’un moteur de table Kafka github_out_queue pointe vers ce topic.github_out_mv qui pointe vers la table GitHub et qui, lorsqu’elle se déclenche, insère des lignes dans le moteur ci-dessus. Les ajouts à la table GitHub seront ainsi envoyés vers notre nouveau topic Kafka.github_out devrait confirmer la bonne réception des messages.Clusters et performance
Utiliser des clusters ClickHouse
Réglage des performances
- Les performances varient selon la taille des messages, le format et les types de tables cibles. Un débit de 100k lignes/s sur un seul table engine peut être considéré comme atteignable. Par défaut, les messages sont lus par blocks, selon le paramètre kafka_max_block_size. Par défaut, celui-ci est défini sur max_insert_block_size, dont la valeur par défaut est 1,048,576. Sauf si les messages sont extrêmement volumineux, cette valeur devrait presque toujours être augmentée. Des valeurs comprises entre 500k et 1M ne sont pas rares. Testez et évaluez l’effet sur le débit.
- Le nombre de consommateurs pour un table engine peut être augmenté avec kafka_num_consumers. Cependant, par défaut, les inserts sont linéarisés dans un seul thread, sauf si kafka_thread_per_consumer est modifié par rapport à sa valeur par défaut de 1. Définissez cette valeur sur 1 pour vous assurer que les flushes sont effectués en parallèle. Notez que créer une table Kafka engine avec N consommateurs (et kafka_thread_per_consumer=1) est logiquement équivalent à créer N Kafka engines, chacun avec une vue matérialisée et kafka_thread_per_consumer=0.
- Augmenter le nombre de consommateurs n’est pas sans coût. Chaque consommateur maintient ses propres buffers et threads, ce qui augmente la surcharge sur le server. Tenez compte de cette surcharge et privilégiez d’abord une mise à l’échelle linéaire sur votre cluster, si possible.
- Si le débit des messages Kafka est variable et que les retards sont acceptables, envisagez d’augmenter
stream_flush_interval_msafin que des blocks plus volumineux soient flushed. - background_message_broker_schedule_pool_size définit le nombre de threads exécutant les tâches d’arrière-plan. Ces threads sont utilisés pour le streaming Kafka. Ce setting est appliqué au démarrage du ClickHouse server et ne peut pas être modifié dans une session utilisateur. Sa valeur par défaut est de 16. Si vous observez des timeouts dans les logs, il peut être judicieux d’augmenter cette valeur.
- Pour la communication avec Kafka, la bibliothèque librdkafka est utilisée, et elle crée elle-même des threads. Un grand nombre de tables Kafka, ou de consommateurs, peut donc entraîner un grand nombre de commutations de contexte. Répartissez cette charge sur le cluster, en ne répliquant que les tables cibles si possible, ou envisagez d’utiliser un table engine pour lire depuis plusieurs topics : une liste de values est prise en charge. Plusieurs vues matérialisées peuvent lire à partir d’une seule table, chacune filtrant les données d’un topic spécifique.
Paramètres supplémentaires
- Kafka_max_wait_ms - Délai d’attente, en millisecondes, avant une nouvelle tentative de lecture des messages depuis Kafka. Ce paramètre est défini au niveau du profil utilisateur et sa valeur par défaut est 5000.