Skip to main content
Ce moteur permet d’intégrer ClickHouse à RabbitMQ. RabbitMQ vous permet de :
  • publier ou de vous abonner à des flux de données ;
  • traiter les flux dès qu’ils sont disponibles.

Créer une table

Paramètres requis :
  • rabbitmq_host_port – hôte:port (par exemple, localhost:5672).
  • rabbitmq_exchange_name – nom de l’exchange RabbitMQ.
  • rabbitmq_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 :
  • rabbitmq_exchange_type – Le type d’exchange RabbitMQ : direct, fanout, topic, headers, consistent_hash. Par défaut : fanout.
  • rabbitmq_routing_key_list – Une liste de clés de routage séparées par des virgules.
  • rabbitmq_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 ainsi que le nom de l’objet racine schema.capnp:Message.
  • rabbitmq_num_consumers – Le nombre de consommateurs par table. Spécifiez davantage de consommateurs si le débit d’un consommateur est insuffisant. Par défaut : 1
  • rabbitmq_num_queues – Nombre total de files. Augmenter ce nombre peut améliorer considérablement les performances. Par défaut : 1.
  • rabbitmq_queue_base - Spécifie un préfixe pour les noms de file. Les cas d’utilisation de ce paramètre sont décrits ci-dessous.
  • rabbitmq_persistent - Si défini sur 1 (true), le mode de remise de la requête d’insertion sera défini sur 2 (marque les messages comme ‘persistent’). Par défaut : 0.
  • rabbitmq_skip_broken_messages – Tolérance de l’analyseur de messages RabbitMQ aux messages incompatibles avec le schéma par bloc. Si rabbitmq_skip_broken_messages = N, le moteur ignore alors N messages RabbitMQ qui ne peuvent pas être analysés (un message équivaut à une ligne de données). Par défaut : 0.
  • rabbitmq_max_block_size - Nombre de lignes collectées avant le vidage des données depuis RabbitMQ. Par défaut : max_insert_block_size.
  • rabbitmq_flush_interval_ms - Délai d’attente avant le vidage des données depuis RabbitMQ. Par défaut : stream_flush_interval_ms.
  • rabbitmq_queue_settings_list - permet de définir les paramètres RabbitMQ lors de la création d’une file. Paramètres disponibles : x-max-length, x-max-length-bytes, x-message-ttl, x-expires, x-priority, x-max-priority, x-overflow, x-dead-letter-exchange, x-queue-type. Le paramètre durable est activé automatiquement pour la file.
  • rabbitmq_address - Adresse de connexion : amqp(s)://user:password@host:port/vhost. Utilisez soit ce paramètre, soit rabbitmq_host_port ; si les deux sont définis, rabbitmq_address est utilisé. Son hôte et son port sont vérifiés par rapport à remote_url_allow_hosts.
  • rabbitmq_vhost - vhost RabbitMQ. Par défaut : '/'.
  • rabbitmq_queue_consume - Utilise des files définies par l’utilisateur et n’effectue aucune configuration RabbitMQ : déclaration des exchanges, des files et des liaisons. Par défaut : false.
  • rabbitmq_username - Nom d’utilisateur RabbitMQ.
  • rabbitmq_password - Mot de passe RabbitMQ.
  • reject_unhandled_messages - Rejette les messages (envoie un accusé de réception négatif RabbitMQ) en cas d’erreur. Ce paramètre est activé automatiquement si un x-dead-letter-exchange est défini dans rabbitmq_queue_settings_list.
  • rabbitmq_commit_on_select - Valide les messages lorsqu’une requête SELECT est effectuée. Par défaut : false.
  • rabbitmq_max_rows_per_message — Le nombre maximal de lignes écrites dans un message RabbitMQ pour les formats basés sur les lignes. Par défaut : 1.
  • rabbitmq_empty_queue_backoff_start_ms — Point de départ du backoff pour replanifier la lecture si la file RabbitMQ est vide.
  • rabbitmq_empty_queue_backoff_end_ms — Point de fin du backoff pour replanifier la lecture si la file RabbitMQ est vide.
  • rabbitmq_empty_queue_backoff_step_ms — Pas du backoff pour replanifier la lecture si la file RabbitMQ est vide.
  • rabbitmq_handle_error_mode — Mode de gestion des erreurs du moteur RabbitMQ. Valeurs possibles : default (une exception sera levée si l’analyse d’un message échoue), stream (le message d’exception et le message brut seront enregistrés dans les colonnes virtuelles _error et _raw_message), dead_letter_queue (les données liées aux erreurs seront enregistrées dans system.dead_letter_queue).

Connexion SSL

Avec la forme rabbitmq_host_port, définissez rabbitmq_secure = 1 pour utiliser TLS. Avec la forme rabbitmq_address, le transport est déterminé par le schéma URI ; utilisez donc amqps : rabbitmq_address = 'amqps://guest:guest@localhost/vhost'. rabbitmq_secure est ignoré pour cette forme d’adresse, et rabbitmq_secure = 1 associé à une adresse amqp:// en texte brut est rejeté plutôt que d’établir silencieusement une connexion en clair. Par défaut, la bibliothèque utilisée ne vérifie pas si la connexion TLS établie est suffisamment sécurisée. Que le certificat soit expiré, autosigné, absent ou invalide, la connexion est simplement autorisée. Une vérification plus stricte des certificats pourra être mise en œuvre à l’avenir. Des paramètres de format peuvent également être ajoutés avec les paramètres liés à RabbitMQ. Exemple :
La configuration du serveur RabbitMQ doit être ajoutée dans le fichier de configuration de ClickHouse. Configuration requise :
Configuration supplémentaire :

Description

SELECT n’est pas particulièrement utile pour lire des 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 RabbitMQ et considérez-le 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 place dans une table créée précédemment.
Lorsque la MATERIALIZED VIEW est liée au moteur, elle commence à collecter les données en arrière-plan. Cela vous permet de recevoir en continu des messages depuis RabbitMQ et de les convertir au format requis à l’aide de SELECT. Une table RabbitMQ peut avoir autant de vues matérialisées que vous le souhaitez. Les données peuvent être acheminées en fonction de rabbitmq_exchange_type et de la rabbitmq_routing_key_list spécifiée. Il ne peut y avoir qu’un seul exchange par table. Un exchange peut être partagé entre plusieurs tables, ce qui permet un routage vers plusieurs tables en même temps. Options de type d’exchange :
  • direct - Le routage est basé sur la correspondance exacte des clés. Exemple de liste de clés de table : key1,key2,key3,key4,key5 ; la clé du message peut être égale à l’une d’elles.
  • fanout - Routage vers toutes les tables (où le nom de l’exchange est identique), quelles que soient les clés.
  • topic - Le routage est basé sur des patterns avec des clés séparées par des points. Exemples : *.logs, records.*.*.2020, *.2018,*.2019,*.2020.
  • headers - Le routage est basé sur des correspondances key=value avec un paramètre x-match=all ou x-match=any. Exemple de liste de clés de table : x-match=all,format=logs,type=report,year=2020.
  • consistent_hash - Les données sont réparties uniformément entre toutes les tables liées (où le nom de l’exchange est identique). Notez que ce type d’exchange doit être activé avec le plugin RabbitMQ : rabbitmq-plugins enable rabbitmq_consistent_hash_exchange.
Le paramètre rabbitmq_queue_base peut être utilisé dans les cas suivants :
  • permettre à différentes tables de partager des files, afin que plusieurs consommateurs puissent être enregistrés pour les mêmes files, ce qui améliore les performances. Si vous utilisez les paramètres rabbitmq_num_consumers et/ou rabbitmq_num_queues, une correspondance exacte des files est obtenue si ces paramètres sont identiques.
  • pouvoir reprendre la lecture à partir de certaines files durables lorsque tous les messages n’ont pas été consommés avec succès. Pour reprendre la consommation à partir d’une file spécifique, définissez son nom dans le paramètre rabbitmq_queue_base et ne spécifiez pas rabbitmq_num_consumers ni rabbitmq_num_queues (valeur par défaut : 1). Pour reprendre la consommation à partir de toutes les files déclarées pour une table spécifique, indiquez simplement les mêmes paramètres : rabbitmq_queue_base, rabbitmq_num_consumers, rabbitmq_num_queues. Par défaut, les noms des files seront uniques pour chaque table.
  • réutiliser des files, puisqu’elles sont déclarées durables et ne sont pas supprimées automatiquement. (Elles peuvent être supprimées via n’importe lequel des outils CLI de RabbitMQ.)
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é formé dans les stream_flush_interval_ms millisecondes, les données seront écrites dans la table, même si le bloc n’est pas complet. Si les paramètres rabbitmq_num_consumers et/ou rabbitmq_num_queues sont spécifiés avec rabbitmq_exchange_type, alors :
  • le plugin rabbitmq-consistent-hash-exchange doit être activé.
  • la propriété message_id des messages publiés doit être spécifiée (unique pour chaque message/lot).
Pour la requête d’insertion, des métadonnées de message sont ajoutées pour chaque message publié : messageID et l’indicateur republished (true s’il a été publié plus d’une fois) - ils sont accessibles via les en-têtes du message. N’utilisez pas la même table pour les inserts et les vues matérialisées. Exemple :

Colonnes virtuelles

  • _exchange_name - Nom de l’exchange RabbitMQ. Type de données : String.
  • _channel_id - ChannelID sur lequel le consommateur ayant reçu le message a été déclaré. Type de données : String.
  • _delivery_tag - DeliveryTag du message reçu. Limité à ce canal. Type de données : UInt64.
  • _redelivered - indicateur redelivered du message. Type de données : UInt8.
  • _message_id - messageID du message reçu ; non vide s’il a été défini lors de la publication du message. Type de données : String.
  • _timestamp - horodatage du message reçu ; non vide s’il a été défini lors de la publication du message. Type de données : UInt64.
Colonnes virtuelles supplémentaires lorsque rabbitmq_handle_error_mode='stream' :
  • _raw_message - Message brut qui n’a pas pu être analysé correctement. 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é correctement.

Mises en garde

Même si vous pouvez spécifier des expressions de colonne par défaut (comme DEFAULT, MATERIALIZED, ALIAS) dans la définition de la table, elles seront ignorées. À la place, les colonnes seront renseignées avec les valeurs par défaut correspondant à leur type.

Prise en charge des formats de données

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

Durabilité des données en cas de coupure d’alimentation

Le moteur RabbitMQ peut perdre silencieusement des lignes déjà consommées si le cache de pages du système d’exploitation est perdu avant que les données insérées ne soient écrites sur le disque. Après l’envoi d’un lot vers les vues matérialisées dépendantes, le consommateur envoie basic.ack au broker, ce qui lui permet de supprimer ces messages. Toutefois, les lignes insérées ne deviennent durables qu’une fois la part cible synchronisée sur le disque, 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 sur le disque de la part cible, le broker a déjà supprimé les messages et le consommateur les ignore lors de la reconnexion : les lignes sont donc perdues sans erreur et count() est 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 le disque. En revanche, la 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é via les 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 rend les parts insérées durables 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 lot est inséré, y compris les cibles des 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 des 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 <rabbitmq_table> direct avec rabbitmq_commit_on_select = 1, où les messages sont accusés de réception lorsque la lecture arrive à son terme, plutôt qu’après que la destination a écrit une part durable.
Dernière modification le 14 août 2026