Skip to main content
该引擎支持将 ClickHouse 与 RabbitMQ 集成。 RabbitMQ 支持您:
  • 发布或订阅数据流。
  • 在流可用后立即进行处理。

创建表

必需参数:
  • rabbitmq_host_port – host:port (例如 localhost:5672) 。
  • rabbitmq_exchange_name – RabbitMQ 的 exchange 名称。
  • rabbitmq_format – 消息格式。使用与 SQL FORMAT 函数相同的记法,例如 JSONEachRow。更多信息,请参见 Formats 章节。
可选参数:
  • rabbitmq_exchange_type – RabbitMQ exchange 的类型:directfanouttopicheadersconsistent_hash。默认值:fanout
  • rabbitmq_routing_key_list – 以逗号分隔的路由键列表。
  • rabbitmq_schema – 如果 格式 需要 schema 定义,则必须使用此参数。例如,Cap’n Proto 需要提供 schema 文件的路径以及根对象 schema.capnp:Message 的名称。
  • rabbitmq_num_consumers – 每个表的消费者数量。如果单个消费者的吞吐量不足,请指定更多消费者。默认值:1
  • rabbitmq_num_queues – 队列总数。增加此数量可显著提升性能。默认值:1
  • rabbitmq_queue_base - 为队列名称指定提示信息。此设置的使用场景见下文。
  • rabbitmq_persistent - 如果设置为 1 (true),则插入查询的投递模式将设置为 2 (将消息标记为“持久”) 。默认值:0
  • rabbitmq_skip_broken_messages – RabbitMQ 消息解析器对每个块中与 schema 不兼容消息的容忍数量。如果 rabbitmq_skip_broken_messages = N,则引擎会跳过 N 条无法解析的 RabbitMQ 消息 (1 条消息等于 1 行数据) 。默认值:0
  • rabbitmq_max_block_size - 从 RabbitMQ 刷新数据前收集的行数。默认值:max_insert_block_size
  • rabbitmq_flush_interval_ms - 从 RabbitMQ 刷新数据的超时时间。默认值:stream_flush_interval_ms
  • rabbitmq_queue_settings_list - 允许在创建队列时设置 RabbitMQ 参数。可用设置:x-max-lengthx-max-length-bytesx-message-ttlx-expiresx-priorityx-max-priorityx-overflowx-dead-letter-exchangex-queue-type。队列的 durable 设置会自动启用。
  • rabbitmq_address - 连接地址:amqp(s)://user:password@host:port/vhost。使用此设置或 rabbitmq_host_port;如果两者都已设置,则使用 rabbitmq_address。其主机和端口会根据 remote_url_allow_hosts 进行检查。
  • rabbitmq_vhost - RabbitMQ vhost。默认值:'/'
  • rabbitmq_queue_consume - 使用用户定义的队列,且不执行任何 RabbitMQ 设置:声明 exchanges、队列、bindings。默认值:false
  • rabbitmq_username - RabbitMQ 用户名。
  • rabbitmq_password - RabbitMQ 密码。
  • reject_unhandled_messages - 发生错误时拒绝消息 (向 RabbitMQ 发送负确认) 。如果在 rabbitmq_queue_settings_list 中定义了 x-dead-letter-exchange,则会自动启用此设置。
  • rabbitmq_commit_on_select - 在执行 select 查询时提交消息。默认值:false
  • rabbitmq_max_rows_per_message — 对于按行组织的格式,单条 RabbitMQ 消息中写入的最大行数。默认值:1
  • rabbitmq_empty_queue_backoff_start_ms — 当 RabbitMQ 队列为空时,重新调度读取的 backoff 起始点。
  • rabbitmq_empty_queue_backoff_end_ms — 当 RabbitMQ 队列为空时,重新调度读取的 backoff 结束点。
  • rabbitmq_empty_queue_backoff_step_ms — 当 RabbitMQ 队列为空时,重新调度读取的 backoff 步长。
  • rabbitmq_handle_error_mode — RabbitMQ 引擎 的错误处理方式。Possible values: default (如果消息解析失败则抛出 exception) 、stream (exception 消息和原始消息将保存在虚拟列 _error_raw_message 中) 、dead_letter_queue (与错误相关的数据将保存在 system.dead_letter_queue 中) 。

SSL 连接

使用 rabbitmq_host_port 形式时,设置 rabbitmq_secure = 1 以使用 TLS。 使用 rabbitmq_address 形式时,传输方式由 URI 方案决定,因此请使用 amqpsrabbitmq_address = 'amqps://guest:guest@localhost/vhost'。对于地址形式,rabbitmq_secure 会被忽略;同时使用 rabbitmq_secure = 1 和明文 amqp:// 地址会被拒绝,而不会在未加密的情况下静默连接。 所使用库的默认行为是,不会检查所建立的 TLS 连接是否足够安全。无论证书已过期、为自签名、缺失还是无效,都会允许建立连接。未来可能会实现更严格的证书检查。 也可以在添加 RabbitMQ 相关设置的同时,添加格式设置。 示例:
应通过 ClickHouse 配置文件添加 RabbitMQ 服务器配置。 所需配置:
附加配置:

说明

SELECT 并不特别适合读取消息 (调试场景除外) ,因为每条消息只能读取一次。更实用的做法是使用 materialized views 创建实时处理链路。为此:
  1. 使用该引擎创建一个 RabbitMQ 消费者,并将其视为数据 stream。
  2. 创建一个具有所需结构的表。
  3. 创建一个 materialized view,将来自该引擎的数据转换后写入之前创建的表中。
MATERIALIZED VIEW 连接到该引擎时,它会在后台开始收集数据。这样你就可以持续从 RabbitMQ 接收消息,并使用 SELECT 将其转换为所需格式。 一个 RabbitMQ 表可以拥有任意多个 materialized views。 可以根据 rabbitmq_exchange_type 和指定的 rabbitmq_routing_key_list 对数据进行路由。 每个表最多只能有一个 exchange。一个 exchange 可以由多个表共享,这样就能同时路由到多个表。 exchange 类型选项:
  • direct - 路由基于 key 的精确匹配。示例表 key 列表:key1,key2,key3,key4,key5,消息 key 可以等于其中任意一个。
  • fanout - 路由到所有表 (exchange 名称相同的表) ,与 key 无关。
  • topic - 路由基于以点分隔的 key pattern。示例:*.logsrecords.*.*.2020*.2018,*.2019,*.2020
  • headers - 路由基于 key=value 匹配,并使用设置 x-match=allx-match=any。示例表 key 列表:x-match=all,format=logs,type=report,year=2020
  • consistent_hash - 数据会在所有已绑定的表之间均匀分布 (exchange 名称相同的表) 。请注意,此 exchange 类型必须通过 RabbitMQ plugin 启用:rabbitmq-plugins enable rabbitmq_consistent_hash_exchange
设置 rabbitmq_queue_base 可用于以下场景:
  • 让不同的表共享队列,以便为同一组队列注册多个消费者,从而获得更好的性能。如果使用 rabbitmq_num_consumers 和/或 rabbitmq_num_queues 设置,那么在这些参数相同的情况下,可以实现队列的精确匹配。
  • 在并非所有消息都被成功消费时,能够从某些持久化队列恢复读取。要从某个特定队列恢复消费,请在 rabbitmq_queue_base 设置中指定其名称,并且不要指定 rabbitmq_num_consumersrabbitmq_num_queues (默认为 1) 。要从为特定表声明的所有队列恢复消费,只需指定相同的设置:rabbitmq_queue_baserabbitmq_num_consumersrabbitmq_num_queues。默认情况下,队列名称对各表都是唯一的。
  • 复用队列,因为它们被声明为持久化且不会被自动删除。 (可通过任意 RabbitMQ CLI 工具删除。)
为了提高性能,接收到的消息会被分组为大小为 max_insert_block_size 的块。如果该块未能在 stream_flush_interval_ms 毫秒内形成,则无论块是否完整,数据都会被 flush 到表中。 如果在指定 rabbitmq_exchange_type 的同时还指定了 rabbitmq_num_consumers 和/或 rabbitmq_num_queues 设置,那么:
  • 必须启用 rabbitmq-consistent-hash-exchange plugin。
  • 必须指定已发布消息的 message_id 属性 (每条消息/批次唯一) 。
对于插入查询,会为每条已发布消息添加消息元数据:messageIDrepublished 标志 (如果消息被发布超过一次,则为 true) ——可通过消息请求头访问。 不要将同一个表同时用于 inserts 和 materialized views。 示例:

虚拟列

  • _exchange_name - RabbitMQ exchange 名称。数据类型:String
  • _channel_id - 接收该消息的消费者所声明的 ChannelID。数据类型:String
  • _delivery_tag - 已接收消息的 DeliveryTag。每个 channel 内独立。数据类型:UInt64
  • _redelivered - 消息的 redelivered 标志。数据类型:UInt8
  • _message_id - 已接收消息的 messageID;如果在消息发布时设置了该值,则非空。数据类型:String
  • _timestamp - 已接收消息的时间戳;如果在消息发布时设置了该值,则非空。数据类型:UInt64
rabbitmq_handle_error_mode='stream' 时,还会有以下附加虚拟列:
  • _raw_message - 无法成功解析的原始消息。数据类型:Nullable(String)
  • _error - 解析失败时产生的异常消息。数据类型:Nullable(String)
注意:只有在解析过程中发生异常时,才会填充 _raw_message_error 虚拟列;如果消息解析成功,它们始终为 NULL

注意事项

即使你可以在表定义中指定默认列表达式 (例如 DEFAULTMATERIALIZEDALIAS) ,这些设置也会被忽略。相反,各列会填充为其类型对应的默认值。

数据格式支持

RabbitMQ 引擎 支持 ClickHouse 支持的所有格式。 单条 RabbitMQ 消息中的行数取决于格式是按行还是按块组织的:
  • 对于按行组织的格式,可通过设置 rabbitmq_max_rows_per_message 控制单条 RabbitMQ 消息中的行数。
  • 对于按块组织的格式,我们无法将块进一步拆分为更小的部分,但单个块中的行数可以通过通用设置 max_block_size 控制。

断电时的数据持久性

如果已插入的数据尚未写入磁盘,OS page cache 就被丢弃,RabbitMQ 引擎可能会在不报错的情况下丢失已消费的行。批次被推送到依赖的 materialized view 后,消费者会向 broker 发送 basic.ack,broker 随后即可删除这些消息。然而,插入的行只有在目标 parts 完成 fsync 后才真正持久化;默认情况下,这一过程不会同步执行 (fsync_after_insert = 0) 。如果在发送确认后、目标 parts 完成 fsync 前 page cache 丢失,broker 已删除消息,而消费者重新连接后会跳过这些消息继续消费,因此这些行会在没有任何错误的情况下丢失,count() 只会变小。普通的进程 kill 不会暴露此问题,因为 kernel 会保留 page cache,并最终将其写回。page cache 丢失则会暴露此问题,例如设备级断电,以及主机或 kernel 的非正常重置。 对于推荐的 materialized view 消费路径 (仅在整个插入管道完成后才发送确认) ,在目标 MergeTree 表上设置 fsync_after_insert = 1 (以及 fsync_part_directory = 1) ,可确保插入的 parts 在发送确认前持久化,从而显著缩小这一风险窗口。必须在批次写入的每个 MergeTree 表上启用该设置,包括级联 materialized view 的目标表;任何仍使用默认设置的此类表仍可能丢失其 parts。异步中间层不会仅凭此设置获得持久性:例如,当 distributed_foreground_insert = 0 时,Distributed 目标端会在后台插入;这在 ClickHouse Cloud 之外是默认行为,因此还需要为其配置自身的持久性设置,或改用同步插入。此缓解措施也不适用于设置了 rabbitmq_commit_on_select = 1 的直接 INSERT ... SELECT ... FROM <rabbitmq_table>;在这种情况下,消息会在读取结束时得到确认,而非在目标端写入持久 parts 后。
最后修改于 2026年8月14日