Skip to main content
Ce moteur permet d’intégrer ClickHouse à NATS. NATS vous permet de :
  • Publier sur des sujets ou vous y abonner.
  • Traiter les nouveaux messages dès qu’ils sont disponibles.

Création d’une table

Paramètres requis :
  • nats_url – hôte:port (par exemple, localhost:4222)..
  • nats_subjects – Liste des sujets auxquels la table NATS doit s’abonner ou publier. Prend en charge les sujets génériques comme foo.*.bar ou baz.>
  • nats_format – Format du message. Utilise la même notation que la fonction SQL FORMAT, par exemple JSONEachRow. Pour plus d’informations, consultez la section Formats.
Paramètres facultatifs :
  • nats_schema – Paramètre à utiliser si le format nécessite une définition de schéma. Par exemple, Cap’n Proto nécessite le chemin du fichier de schéma et le nom de l’objet racine schema.capnp:Message.
  • nats_stream – Nom d’un stream existant dans NATS JetStream.
  • nats_consumer_name – Nom d’un durable pull consommateur existant dans NATS JetStream.
  • nats_num_consumers – Nombre de consommateurs par table. Par défaut : 1. Spécifiez davantage de consommateurs si le throughput d’un consommateur est insuffisant, uniquement pour NATS core.
  • nats_queue_group – Nom du groupe de queue pour les abonnés NATS. La valeur par défaut est le nom de la table.
  • nats_max_reconnect – Obsolète et sans effet ; la reconnexion est effectuée en permanence avec le délai d’expiration nats_reconnect_wait.
  • nats_reconnect_wait – Durée d’attente, en millisecondes, entre chaque tentative de reconnexion. Par défaut : 2000.
  • nats_server_list - Liste de serveurs pour la connexion. Peut être spécifiée pour se connecter à un cluster NATS.
  • nats_skip_broken_messages - Tolérance de l’analyseur de messages NATS aux messages incompatibles avec le schéma par bloc. Par défaut : 0. Si nats_skip_broken_messages = N, alors le moteur ignore N messages NATS qui ne peuvent pas être analysés (un message équivaut à une ligne de données).
  • nats_max_block_size - Nombre de lignes collectées par poll(s) pour vider les données depuis NATS. Par défaut : max_insert_block_size.
  • nats_flush_interval_ms - Délai d’expiration pour le vidage des données lues depuis NATS. Par défaut : stream_flush_interval_ms.
  • nats_wait_for_flush_interval - Si true, un cycle de streaming en arrière-plan reste ouvert pendant tout l’intervalle de vidage (nats_flush_interval_ms, ou stream_flush_interval_ms sinon) au lieu de se terminer dès que la queue du consommateur est vidée, ce qui permet d’accumuler davantage de messages dans un seul bloc au prix d’une latence d’ingestion supplémentaire pouvant aller jusqu’à un intervalle de vidage. Par défaut : false (comportement de vidage immédiat à faible latence).
  • nats_username - Nom d’utilisateur NATS. Lorsqu’il est stocké dans une collection nommée définie dans le fichier de configuration du serveur, la requête ne peut pas remplacer nats_url ou nats_server_list de la collection.
  • nats_password - Mot de passe NATS. Lorsqu’il est stocké dans une collection nommée définie dans le fichier de configuration du serveur, la requête ne peut pas remplacer nats_url ou nats_server_list de la collection.
  • nats_token - Jeton d’authentification NATS. Lorsqu’il est stocké dans une collection nommée définie dans le fichier de configuration du serveur, la requête ne peut pas remplacer nats_url ou nats_server_list de la collection.
  • nats_credential_file - Chemin vers un fichier d’identifiants NATS. Il est accepté uniquement depuis une collection nommée définie dans le fichier de configuration du serveur, dont nats_url et nats_server_list ne sont pas remplacés par la requête, car le serveur ouvre le chemin avec ses propres privilèges. Dans une requête, transmettez plutôt le contenu du fichier dans nats_credentials.
  • nats_credentials - Contenu des identifiants NATS (la même charge utile que dans un fichier .creds avec le JWT utilisateur et la seed). Comme c’est la seule syntaxe qu’une requête peut utiliser, elle remplace un nats_credential_file hérité d’une collection nommée au lieu d’entrer en conflit avec lui, sauf si l’opérateur a verrouillé ce chemin avec <nats_credential_file overridable="false">. Une chaîne vide ne peut pas lui être attribuée pour supprimer les identifiants fournis par une collection nommée.
  • nats_ca_file - Chemin vers un fichier contenant les certificats CA de confiance utilisés pour vérifier le certificat du serveur NATS. Nécessite nats_secure. Comme nats_credential_file, il est accepté uniquement depuis une collection nommée définie dans le fichier de configuration du serveur, dont nats_url et nats_server_list ne sont pas remplacés par la requête, car le serveur ouvre le chemin avec ses propres privilèges.
  • nats_client_cert_file - Chemin vers le certificat client présenté au serveur NATS. Nécessite nats_secure et nats_client_key_file. Accepté depuis les mêmes sources que nats_ca_file.
  • nats_client_key_file - Chemin vers la clé privée de nats_client_cert_file. Accepté depuis les mêmes sources que nats_ca_file.
  • nats_startup_connect_tries - Nombre de tentatives de connexion au démarrage. Par défaut : 5.
  • nats_max_rows_per_message — Nombre maximal de lignes écrites dans un message NATS pour les formats basés sur les lignes. (par défaut : 1).
  • nats_commit_on_select - Valider les messages lorsqu’une requête est effectuée. S’applique uniquement à JetStream ; NATS core ne comporte pas d’accusés de réception. Par défaut : 0.
  • nats_handle_error_mode — Comment gérer les erreurs pour le moteur NATS. Valeurs possibles : default (une exception sera levée en cas d’échec de l’analyse d’un message), stream (le message d’exception et le message brut seront enregistrés dans les colonnes virtuelles _error et _raw_message).
Connexion SSL : Pour utiliser une connexion sécurisée, définissez nats_secure = 1. La vérification du certificat est contrôlée par la variable d’environnement CLICKHOUSE_NATS_TLS_SECURE ; Si le certificat est expiré, auto-signé, manquant ou non valide pour toute autre raison, désactivez la vérification en définissant CLICKHOUSE_NATS_TLS_SECURE=0. Un certificat de serveur signé par une CA privée est vérifié en indiquant le certificat de CA via nats_ca_file, ce qui est préférable à la désactivation de la vérification. Lorsque le serveur exige des certificats clients, fournissez-les avec nats_client_cert_file et nats_client_key_file. Ces trois éléments sont des paramètres de l’opérateur : ils proviennent d’une collection nommée définie dans le fichier de configuration du serveur. Chaque fichier est lu lorsque la table établit sa connexion ; un fichier illisible ou malformé fait donc échouer la requête plutôt que la négociation de connexion. Écriture dans une table NATS : Si la table lit depuis un seul sujet, tout insert publiera vers ce même sujet. En revanche, si la table lit depuis plusieurs sujets, il faut préciser vers quel sujet nous souhaitons publier. C’est pourquoi, lors de l’insertion dans une table comportant plusieurs sujets, le paramètre stream_like_engine_insert_queue est nécessaire. Vous pouvez sélectionner l’un des sujets depuis lesquels la table lit et y publier vos données. Par exemple :
Des paramètres de format peuvent également être ajoutés en plus des paramètres liés à NATS. Exemple :
Vous pouvez ajouter la configuration du serveur NATS à l’aide du fichier de configuration ClickHouse. Plus précisément, vous pouvez ajouter votre mot de passe pour le moteur NATS :

Description

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 ce faire :
  1. Utilisez le moteur pour créer un consommateur NATS et le considérer comme un flux de données.
  2. Créez une table avec la structure souhaitée.
  3. Créez une vue matérialisée qui convertit les données du moteur et les insère dans une table créée au préalable.
Lorsque la MATERIALIZED VIEW est rattachée au moteur, elle commence à collecter les données en arrière-plan. Cela vous permet de recevoir en continu des messages depuis NATS et de les convertir au format requis à l’aide de SELECT. Une table NATS peut avoir autant de vues matérialisées que vous le souhaitez ; elles ne lisent pas directement les données de la table, mais reçoivent les nouveaux enregistrements (par blocs). Vous pouvez ainsi écrire dans plusieurs tables avec différents niveaux de détail (avec regroupement/agrégation ou sans). Exemple :
Pour ne plus recevoir de données de flux ou pour modifier la logique de conversion, détachez la vue matérialisée :
Si vous souhaitez modifier la table cible en utilisant ALTER, nous vous recommandons de désactiver la vue matérialisée afin d’éviter des écarts entre la table cible et les données de la vue.

Colonnes virtuelles

  • _subject - Sujet du message NATS. Type de données : String.
Colonnes virtuelles supplémentaires lorsque nats_handle_error_mode='stream' :
  • _raw_message - Message brut qui n’a pas pu être analysé avec succès. Type de données : Nullable(String).
  • _error - Message d’exception survenu lors de l’échec de l’analyse. Type de données : Nullable(String).
Remarque : les colonnes virtuelles _raw_message et _error ne sont renseignées qu’en cas d’exception lors de l’analyse ; elles sont toujours NULL lorsque le message a été analysé avec succès.

Prise en charge des formats de données

Le moteur NATS prend en charge tous les formats pris en charge par ClickHouse. Le nombre de lignes dans un message NATS dépend du fait que le format soit basé sur des lignes ou sur des blocs :
  • Pour les formats basés sur des lignes, le nombre de lignes dans un message NATS peut être contrôlé en définissant nats_max_rows_per_message.
  • Pour les formats basés sur des 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.

Utilisation de JetStream

Avant d’utiliser le moteur NATS avec NATS JetStream, vous devez créer un stream NATS et un durable pull consommateur. Pour cela, vous pouvez par exemple utiliser l’utilitaire nats du paquet NATS CLI :
Après avoir créé le stream et le durable pull consommateur, nous pouvons créer une table avec le moteur NATS. Pour cela, vous devez renseigner : nats_stream, nats_consumer_name et nats_subjects:
Les tables JetStream garantissent une livraison au moins une fois : un message n’est acquitté qu’après avoir été inséré dans les vues matérialisées dépendantes. Ainsi, un message dont l’insertion échoue ou est interrompue reste non acquitté et est livré à nouveau. NATS Core (sans JetStream) ne prend en charge ni les accusés de réception ni la relecture ; il garantit donc une livraison au plus une fois, et un message interrompu est perdu.

Durabilité des données

Cette section s’applique uniquement à JetStream. NATS Core n’a pas d’accusés de réception et garantit une livraison au plus une fois, comme décrit ci-dessus. Il n’existe donc aucune fenêtre pendant laquelle un message acquitté peut être perdu. Une table JetStream peut perdre silencieusement des lignes déjà consommées si le cache de pages de l’OS est perdu avant que les données insérées ne soient écrites sur disque. Après l’envoi d’un batch aux vues matérialisées dépendantes, le consommateur accuse réception de ces messages, ce qui permet au stream de progresser au-delà. Les lignes insérées ne sont toutefois durables qu’une fois la part cible synchronisée avec fsync, ce qui ne se produit pas de manière synchrone par défaut (fsync_after_insert = 0). Si le cache de pages est perdu après l’accusé de réception, mais avant la synchronisation de la part cible avec fsync, les messages ne sont plus renvoyés et les lignes sont donc perdues sans erreur, count() étant simplement inférieur. Un simple arrêt forcé du processus ne révèle pas ce problème, car le kernel conserve le cache de pages et finit par l’écrire sur disque. En revanche, une perte du cache de pages le révèle, par exemple en cas de coupure d’alimentation au niveau du périphérique ou de réinitialisation non propre de l’hôte ou du kernel. Pour le chemin de consommation recommandé des vues matérialisées (l’accusé de réception n’est envoyé qu’une fois l’ensemble du pipeline d’insertion terminé), définir fsync_after_insert = 1 (et fsync_part_directory = 1) sur les tables MergeTree cibles garantit la durabilité des parts insérées avant l’envoi de l’accusé de réception, ce qui réduit considérablement cette fenêtre. Le paramètre doit être activé sur chaque table MergeTree dans laquelle le batch est inséré, y compris les cibles de vues matérialisées en cascade ; toute table de ce type laissée avec la valeur par défaut peut encore perdre sa part. Les intermédiaires asynchrones ne gagnent pas en durabilité avec ce seul paramètre : par exemple, une cible Distributed effectue les insertions en arrière-plan lorsque distributed_foreground_insert = 0, ce qui est la valeur par défaut en dehors de ClickHouse Cloud. Elle nécessite donc ses propres paramètres de durabilité ou une insertion synchrone. Cette mesure d’atténuation ne s’applique pas non plus à un INSERT ... SELECT ... FROM <nats_table> direct avec nats_commit_on_select = 1, où les messages sont confirmés lorsque la lecture atteint sa fin plutôt qu’après que la destination a écrit une part durable.
Dernière modification le 26 septembre 2026