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
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = NATS SETTINGS
nats_url = 'host:port',
nats_subjects = 'subject1,subject2,...',
nats_format = 'data_format'[,]
[nats_schema = '',]
[nats_num_consumers = N,]
[nats_queue_group = 'group_name',]
[nats_secure = false,]
[nats_max_reconnect = N,]
[nats_reconnect_wait = N,]
[nats_server_list = 'host1:port1,host2:port2,...',]
[nats_skip_broken_messages = N,]
[nats_max_block_size = N,]
[nats_flush_interval_ms = N,]
[nats_username = 'user',]
[nats_password = 'password',]
[nats_token = 'clickhouse',]
[nats_credentials = '-----BEGIN NATS USER JWT----- ...',]
[nats_startup_connect_tries = 5,]
[nats_max_rows_per_message = 1,]
[nats_commit_on_select = false,]
[nats_handle_error_mode = 'default']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 commefoo.*.baroubaz.>nats_format– Format du message. Utilise la même notation que la fonction SQLFORMAT, par exempleJSONEachRow. 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 racineschema.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’expirationnats_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. Sinats_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- Sitrue, un cycle de streaming en arrière-plan reste ouvert pendant tout l’intervalle de vidage (nats_flush_interval_ms, oustream_flush_interval_mssinon) 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 remplacernats_urlounats_server_listde 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 remplacernats_urlounats_server_listde 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 remplacernats_urlounats_server_listde 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, dontnats_urletnats_server_listne 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 dansnats_credentials.nats_credentials- Contenu des identifiants NATS (la même charge utile que dans un fichier.credsavec le JWT utilisateur et la seed). Comme c’est la seule syntaxe qu’une requête peut utiliser, elle remplace unnats_credential_filehé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_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_erroret_raw_message).
Connexion SSL :
Pour une connexion sécurisée, utilisez 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 invalide pour une autre raison, désactivez la vérification en définissant CLICKHOUSE_NATS_TLS_SECURE=0.
É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 :
CREATE TABLE queue (
key UInt64,
value UInt64
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1,subject2',
nats_format = 'JSONEachRow';
INSERT INTO queue
SETTINGS stream_like_engine_insert_queue = 'subject2'
VALUES (1, 1);Des paramètres de format peuvent également être ajoutés en plus des paramètres liés à NATS.
Exemple :
CREATE TABLE queue (
key UInt64,
value UInt64,
date DateTime
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1',
nats_format = 'JSONEachRow',
date_time_input_format = 'best_effort';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 :
<nats>
<user>click</user>
<password>house</password>
<token>clickhouse</token>
</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 :
- Utilisez le moteur pour créer un consommateur NATS et le considérer 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 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 :
CREATE TABLE queue (
key UInt64,
value UInt64
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1',
nats_format = 'JSONEachRow',
date_time_input_format = 'best_effort';
CREATE TABLE daily (key UInt64, value UInt64)
ENGINE = MergeTree() ORDER BY key;
CREATE MATERIALIZED VIEW consumer TO daily
AS SELECT key, value FROM queue;
SELECT key, value FROM daily ORDER BY key;Pour ne plus recevoir de données de flux ou pour 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 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 :
création d’un stream
$ nats stream add
? Stream Name stream_name
? Subjects stream_subject
? Storage file
? Replication 1
? Retention Policy Limits
? Discard Policy Old
? Stream Messages Limit -1
? Per Subject Messages Limit -1
? Total Stream Size -1
? Message TTL -1
? Max Message Size -1
? Duplicate tracking time window 2m0s
? Allow message Roll-ups No
? Allow message deletion Yes
? Allow purging subjects or the entire stream Yes
Stream stream_name was created
Information for Stream stream_name created 2025-10-03 14:12:51
Subjects: stream_subject
Replicas: 1
Storage: File
Options:
Retention: Limits
Acknowledgments: true
Discard Policy: Old
Duplicate Window: 2m0s
Direct Get: true
Allows Msg Delete: true
Allows Purge: true
Allows Per-Message TTL: false
Allows Rollups: false
Limits:
Maximum Messages: unlimited
Maximum Per Subject: unlimited
Maximum Bytes: unlimited
Maximum Age: unlimited
Maximum Message Size: unlimited
Maximum Consumers: unlimited
State:
Messages: 0
Bytes: 0 B
First Sequence: 0
Last Sequence: 0
Active Consumers: 0création d’un durable pull consommateur
$ nats consumer add
? Select a Stream stream_name
? Consumer name consumer_name
? Delivery target (empty for Pull Consumers)
? Start policy (all, new, last, subject, 1h, msg sequence) all
? Acknowledgment policy explicit
? Replay policy instant
? Filter Stream by subjects (blank for all)
? Maximum Allowed Deliveries -1
? Maximum Acknowledgments Pending 0
? Deliver headers only without bodies No
? Add a Retry Backoff Policy No
Information for Consumer stream_name > consumer_name created 2025-10-03T14:13:51+03:00
Configuration:
Name: consumer_name
Pull Mode: true
Deliver Policy: All
Ack Policy: Explicit
Ack Wait: 30.00s
Replay Policy: Instant
Max Ack Pending: 1,000
Max Waiting Pulls: 512
State:
Last Delivered Message: Consumer sequence: 0 Stream sequence: 0
Acknowledgment Floor: Consumer sequence: 0 Stream sequence: 0
Outstanding Acks: 0 out of maximum 1,000
Redelivered Messages: 0
Unprocessed Messages: 0
Waiting Pulls: 0 of maximum 512Aprè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:
CREATE TABLE nats_jet_stream (
key UInt64,
value UInt64
) ENGINE NATS
SETTINGS nats_url = 'localhost:4222',
nats_stream = 'stream_name',
nats_consumer_name = 'consumer_name',
nats_subjects = 'stream_subject',
nats_format = 'JSONEachRow';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.