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
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.
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_ca_file- Chemin vers un fichier contenant les certificats CA de confiance utilisés pour vérifier le certificat du serveur NATS. Nécessitenats_secure. Commenats_credential_file, 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.nats_client_cert_file- Chemin vers le certificat client présenté au serveur NATS. Nécessitenats_secureetnats_client_key_file. Accepté depuis les mêmes sources quenats_ca_file.nats_client_key_file- Chemin vers la clé privée denats_client_cert_file. Accepté depuis les mêmes sources quenats_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_erroret_raw_message).
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 :
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.
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 :
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.
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).
_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’utilitairenats du paquet NATS CLI :
création d’un stream
création d’un stream
création d’un durable pull consommateur
création d’un durable pull consommateur
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 avecfsync, 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.