使用原生 MQTT 为数百万客户端提供服务
RabbitMQ 的核心协议一直是 AMQP 0.9.1。为了支持 MQTT、STOMP 和 AMQP 1.0,代理通过其核心协议进行透明代理。虽然这是扩展 RabbitMQ 以支持更多消息协议的简单方法,但它会降低可伸缩性和性能。
在过去 9 个月里,我们重写了 MQTT 插件,使其不再通过 AMQP 0.9.1 进行代理。取而代之的是,MQTT 插件会解析 MQTT 消息并直接将其发送到队列。这就是我们称之为原生 MQTT 的方式。
结果令人瞩目
- 内存使用量下降高达 95%,在数百万连接的情况下节省数百 GB 内存。
- RabbitMQ 首次能够处理数百万个连接。
- 端到端延迟下降 50% - 70%。
- 吞吐量提高 30% - 40%。
原生 MQTT 将 RabbitMQ 转变为一个 MQTT 代理,为更广泛的物联网用例打开了大门。
原生 MQTT 将在 RabbitMQ 3.12 中发布。
概述
如图 1 所示,在 RabbitMQ 3.11 及之前版本中,MQTT 插件的工作方式是解析 MQTT 消息,并通过 AMQP 0.9.1 协议将其转发给通道 (channel),进而将消息路由到队列 (queues)。图 1 中的每个蓝点代表一个 Erlang 进程。对于每个传入的 MQTT 连接,总共会创建 22 个 Erlang 进程。
MQTT 插件(即 Erlang 应用程序 rabbitmq_mqtt)为每个 MQTT 连接创建 16 个 Erlang 进程,其中包括一个负责 MQTT Keep Alive(保持连接) 的进程和一组充当 AMQP 0.9.1 客户端的进程。
Erlang 应用程序 rabbit 为每个 MQTT 连接创建 6 个 Erlang 进程。它们实现了核心 AMQP 0.9.1 服务器中针对每个客户端的部分。
图 2 显示,原生 MQTT 在每个传入的 MQTT 连接上仅需一个 Erlang 进程。
这单个 Erlang 进程负责解析 MQTT 消息、遵守 MQTT Keep Alive、执行身份验证和授权,并将消息路由到队列。
通常的 MQTT 工作负载由许多物联网 (IoT) 设备组成,它们定期向 MQTT 代理发送数据。例如,可能有成千上万甚至数百万个设备,每个设备每隔几秒或几分钟发送一次状态更新。
我们在去年 RabbitMQ 峰会的演讲《RabbitMQ 性能改进》中了解到,创建 Erlang 进程的开销很小。(Erlang 进程比 Java 线程轻量得多。Erlang 进程可以与 Golang 中的 Goroutine 相媲美。)然而,对于 100 万个传入的 MQTT 客户端连接,RabbitMQ 是创建 2200 万个 Erlang 进程(RabbitMQ 3.11 及之前)还是仅创建 100 万个 Erlang 进程(RabbitMQ 3.12),在可扩展性方面有着巨大的差异。这种差异将在内存使用和延迟与吞吐量部分进行分析。
不仅 rabbitmq_mqtt 插件被重写,rabbitmq_web_mqtt 插件也进行了重写。这意味着从 3.12 版本开始,每个传入的 MQTT over WebSocket 连接也将只有一个 Erlang 进程。因此,本博客文章中概述的所有性能改进同样适用于 MQTT over WebSockets。
“原生 MQTT”并非 MQTT 规范中的官方术语。“原生 MQTT”指的是新的 RabbitMQ 3.12 实现,其中 RabbitMQ “原生”支持 MQTT,即 MQTT 流量不再通过 AMQP 0.9.1 进行代理。
新的 MQTT QoS 0 队列类型
原生 MQTT 附带了一种名为 rabbit_mqtt_qos0_queue 的新型 RabbitMQ 队列。要让 MQTT 插件创建该类型的新队列,必须启用同名的 3.12 功能标志 (feature flag) rabbit_mqtt_qos0_queue。(请记住,功能标志并非旨在用作集群配置的形式。在成功的滚动升级后,您应该启用所有功能标志。每个功能标志将在未来的 RabbitMQ 版本中变为强制性。)
在解释新队列类型之前,我们应该先了解队列与 MQTT 订阅者之间的关系。
在所有 RabbitMQ 版本中,MQTT 插件都会为每个 MQTT 订阅者创建一个专用队列。更准确地说,每个 MQTT 连接可能有 0 个、1 个或 2 个队列。
- 如果 MQTT 客户端从不发送 SUBSCRIBE 数据包,则该 MQTT 连接没有队列。此时 MQTT 客户端仅发布消息。
- 如果 MQTT 客户端创建了一个或多个具有相同服务质量 (QoS) 级别的订阅,则该 MQTT 连接有 1 个队列。
- 如果 MQTT 客户端创建了一个或多个同时包含 QoS 0(最多一次) 和 QoS 1(至少一次) 的订阅,则该 MQTT 连接有 2 个队列。
当列出队列时,您会观察到 mqtt-subscription-<MQTT client ID>qos[0|1] 这种队列命名模式,其中 <MQTT client ID> 是 MQTT 客户端标识符,而 [0|1] 对于 QoS 0 订阅为 0,对于 QoS 1 订阅为 1。为每个 MQTT 订阅者提供单独的队列是有意义的,因为每个 MQTT 订阅者都会收到其专属的应用程序消息副本。
默认情况下,MQTT 插件会创建经典队列 (classic queues)。
MQTT 插件会为订阅 MQTT 的客户端透明地创建队列。MQTT 规范并未定义队列的概念,MQTT 客户端也不知道这些队列的存在。队列是 RabbitMQ 实现 MQTT 协议的一种实现细节。
图 3 显示了一个以 CleanSession=1 连接并以 QoS 0 订阅的 MQTT 订阅者。“清理会话” (Clean session) 意味着 MQTT 会话仅在客户端和服务器之间的网络连接存在期间有效。当会话结束时,服务器中的所有会话状态都将被删除,这意味着队列会被自动删除。
如图 3 所示,每个经典队列都会导致两个额外的 Erlang 进程:一个 监督 (supervisor) 进程和一个工作 (worker) 进程。
新的队列类型优化工作原理如下:如果
- 功能标志
rabbit_mqtt_qos0_queue已启用,且 - MQTT 客户端以
CleanSession=1连接,且 - MQTT 客户端以 QoS 0 进行订阅,
那么 MQTT 插件将创建一个 rabbit_mqtt_qos0_queue 类型的队列,而不是经典队列。
图 4 显示,这种新的队列类型是一种“伪”队列或“虚拟”队列:它与您熟知的队列类型(经典队列、仲裁队列和流)非常不同,因为它既不是一个单独的 Erlang 进程,也不会将消息存储在磁盘上。相反,这种队列类型是 Erlang 进程邮箱的一个子集。MQTT 消息直接发送到订阅客户端的 MQTT 连接进程。换句话说,MQTT 消息直接发送给所有“在线”的 MQTT 订阅者。
将其理解为跳过队列(正如这篇 RabbitMQ 2022 年峰会演讲幻灯片最后一行 no queue at all? 🤔 所指)会更准确。将消息直接发送到 MQTT 连接进程实现为一种队列类型,是为了简化消息路由和协议互操作性,使得消息不仅可以从 MQTT 发布连接进程发送,也可以从通道进程发送。后者使得从 AMQP 0.9.1、AMQP 1.0 或 STOMP 客户端可以直接向 MQTT 订阅者连接进程发送消息,从而跳过专门的队列进程。
我们现在了解到这种新的队列类型跳过了队列进程。然而,在 MQTT 环境中这样做有什么优势呢?MQTT 工作负载的特征是有许多设备向代理发送数据并从中接收数据。以下是该新队列类型提供优化的四个原因,按重要性递减排列:
优势 1:大规模扇出
实现将消息从“云”(MQTT 代理)大规模扇出发送到所有设备。
对于经典队列和仲裁队列,每个队列客户端(即通道进程或 MQTT 连接进程)都会为所有目标队列保持状态,以进行流控。例如,通道进程会在其 进程字典 (process dictionary) 中保存来自每个目标队列的信用额度。(阅读我们关于基于信用的流控的博客文章以了解更多信息。)
- 如果有一个通道进程向 300 万个 MQTT 设备发送消息,那么通道进程字典中会保存 300 万个条目(数百 MB 的内存)。
- 如果有 100 个通道进程,每个都向 300 万个设备发送消息,那么进程字典中总共会有 3 亿个条目(最好不要尝试)。
- 如果有几千个通道或 MQTT 连接进程,每个都向 300 万个设备发送消息,RabbitMQ 将耗尽内存并崩溃。
即使这种大规模扇出极为罕见,例如一天一次,从发送 Erlang 进程到所有目标队列的状态也会一直保存在内存中(直到目标队列被删除)。
新队列类型最重要的特征是其队列客户端是无状态的。这意味着 MQTT 连接或通道进程可以向 300 万个 MQTT 连接进程发送消息(即“即发即弃”,虽然这仍然会暂时消耗大量内存),而无需保持任何目标队列的状态。一旦垃圾回收机制运行,队列客户端进程的内存使用量就会降至 0 MB。
优势 2:更低的内存占用
不仅是在大规模扇出场景,即使在 1:1 拓扑(每个发布者正好发送给一个订阅者)中,新的 rabbit_mqtt_qos0_queue 类型通过跳过队列进程也节省了大量的内存。
即使在 3.12 版本中使用原生 MQTT,如果禁用了 rabbit_mqtt_qos0_queue 功能标志,300 万个以 QoS 0 订阅的 MQTT 设备仍会导致 900 万个 Erlang 进程。启用该标志后,相同的工作负载仅导致 300 万个 Erlang 进程,因为跳过了每个 MQTT 订阅者的额外队列监督进程和队列工作进程,从而节省了数 GB 的进程内存。
优势 3:更低的发布者确认延迟
尽管 MQTT 客户端以 QoS 0 订阅,但另一个 MQTT 客户端仍可以发送 QoS 1 消息(或者等效地,AMQP 0.9.1 发送客户端的通道可以被置于确认模式)。在这种情况下,发布客户端需要来自代理的发布者确认。
新队列类型的客户端(属于 MQTT 发布连接进程或通道进程的一部分)会代表“队列服务器进程”自动进行确认,因为 QoS 0 最多一次消息在从代理到 MQTT 订阅者的途中可能无论如何都会丢失。这带来了更低的发布者确认延迟。
在 RabbitMQ 中,只有当所有目标队列确认已收到消息后,消息才会向发布客户端返回确认。因此,更重要的是,发布过程仅在等待可能具有至少一次消费者的队列的确认。有了新的队列类型,在将一条消息路由到重要的仲裁队列以及 100 万个 MQTT QoS 0 订阅者的场景中,如果 100 万个 MQTT 连接进程中的某一个过载(并因此回复非常缓慢),发布过程不会因等待确认而阻塞。
优势 4:更低的端到端延迟
由于跳过了队列进程,消息传输跳数减少了一次,从而实现了更低的端到端延迟。
过载保护
由于新的队列类型没有流控,MQTT 消息到达 MQTT 连接进程邮箱的速度可能会快于从该进程投递到 MQTT 订阅客户端的速度。这种情况可能在 MQTT 订阅客户端与 RabbitMQ 之间的网络连接较差,或者在许多发布者使单个 MQTT 订阅客户端过载的大规模扇入 (fan-in) 场景中发生。
为了防止因 MQTT QoS 0 消息堆积在 MQTT 连接进程邮箱中而导致高内存使用,RabbitMQ 会在满足以下两个条件时,有意从 rabbit_mqtt_qos0_queue 中丢弃 QoS 0 消息:
- MQTT 连接进程邮箱中的消息数量超过了设置
mqtt.mailbox_soft_limit(默认为 200),并且 - 发送到 MQTT 客户端的套接字处于繁忙状态(发送速度不够快)。
请注意,进程邮箱中可能存在其他消息(例如从 MQTT 订阅客户端发送到 RabbitMQ 的应用程序消息,或来自其他队列类型的确认消息),这些消息显然不会被丢弃。但是,这些其他消息也会计入 mqtt.mailbox_soft_limit。
将 mqtt.mailbox_soft_limit 设置为 0 将禁用过载保护机制,这意味着 QoS 0 消息永远不会被 RabbitMQ 有意丢弃。将该值设置为非常大的值会降低有意丢弃 QoS 0 消息的可能性,同时增加导致集群范围内存告警的风险(特别是当消息负载很大或者存在许多过载的 rabbit_mqtt_qos0_queue 类型队列时)。
mqtt.mailbox_soft_limit 可以被看作是一种队列长度限制,尽管并不精确,因为如前所述,Erlang 进程邮箱中可能包含除 MQTT 应用程序消息之外的其他消息。这就是为什么配置键 mqtt.mailbox_soft_limit 包含单词 soft(软)的原因。所描述的过载保护机制大致对应于您在经典队列和仲裁队列中熟知的 drop-head(丢弃头部)溢出行为。
队列类型命名
不要依赖此新队列类型的名称。功能标志 rabbit_mqtt_qos0_queue 的名称不会改变。然而,如果我们在未来决定在其他 RabbitMQ 用例(例如直接回复 (Direct Reply-to))中重用其部分设计,我们可能会更改 rabbit_mqtt_qos0_queue 队列类型的名称。
事实上,MQTT 客户端应用程序和 RabbitMQ 核心(Erlang 应用程序 rabbit)都不知道这种新队列类型的存在。最终用户也不会察觉到这种新队列类型,除非在管理 UI 中或通过 rabbitmqctl 列出队列时。
既然我们已经理解了原生 MQTT 的架构和 MQTT QoS 0 队列类型的意图,我们可以继续进行性能基准测试。
内存使用
本节比较了 RabbitMQ 3.11 中的 MQTT 和 RabbitMQ 3.12 中的原生 MQTT 的内存使用情况。完整的测试设置可以在 ansd/rabbitmq-mqtt 中找到。
本节中的两个测试是在具有以下 RabbitMQ 配置子集的 3 节点集群上完成的:
mqtt.tcp_listen_options.sndbuf = 1024
mqtt.tcp_listen_options.recbuf = 1024
mqtt.tcp_listen_options.buffer = 1024
management_agent.disable_metrics_collector = true
TCP 缓冲区大小配置得很小,这样它们就不会导致高二进制内存使用。
rabbitmq_management 插件中禁用了指标收集。对于生产环境,应该使用 Prometheus。rabbitmq_management 插件并非旨在处理大量的统计信息发射对象,如队列和连接。
100 万个 MQTT 连接
第一个测试使用 1,000,000 个 MQTT 连接,这些连接在建立连接后仅发送 MQTT Keep Alive。没有发送或接收任何 MQTT 应用程序消息。
如图 5 所示,3 节点集群在 3.11 中需要 108.0 + 100.7 + 92.4 = 301.1 GiB 内存,而在 3.12 中仅需要 6.1 + 6.3 + 6.3 = 18.7 GiB 内存。因此,3.11 比 3.12 需要多 16 倍(或 282 GiB)的内存。原生 MQTT 将内存使用量降低了 94%。

3.11 中导致内存使用的最大部分是进程内存。正如概述部分所解释的,3.12 中的原生 MQTT 每个 MQTT 连接使用 1 个 Erlang 进程,而 3.11 每个 MQTT 连接使用 22 个 Erlang 进程。3.11 中的某些进程会保持大量状态,从而导致高内存使用。
原生 MQTT 的低内存使用不仅通过使用单个 Erlang 进程实现,还通过减少该单个 Erlang 进程的状态来实现。在 PR #5895 中实施了众多的内存优化,例如从进程状态中删除了长列表和不必要的函数引用。
在开发环境上的进一步测试显示,对于总共 900 万个 MQTT 连接,原生 MQTT 在 3 节点集群中每个节点需要约 56 GB 的内存。
10 万个发布者,10 万个订阅者
第二个测试使用了 100,000 个发布者和 100,000 个订阅者。它们以 1:1 的拓扑结构进行发送和接收,这意味着每个发布者每 2 分钟向正好一个 QoS 0 订阅者发送一条 QoS 0 应用程序消息。
如图 6 所示,3 节点集群在 3.11 中需要 21.6 + 21.5 + 21.7 = 64.8 GiB 内存,而在 3.12 中仅需要 2.6 + 2.6 + 2.6 = 7.8 GiB 内存。

在开发环境上的进一步测试显示,对于以下场景,原生 MQTT 在 3 节点集群中每个节点需要约 47 GB 的内存:
- 总共 300 万个 MQTT 连接
- 1:1 拓扑
- 150 万个发布者每 3 分钟发送一次 64 字节负载的 QoS 1 消息
- 150 万个 QoS 1 订阅者(即 150 万个第 2 版经典队列)
RabbitMQ 3.11.6 中发布的 PR #6684 大幅减少了许多经典队列的内存使用,这有利于许多 MQTT QoS 1 或 CleanSession=0 订阅者的用例。
延迟和吞吐量
为了比较 3.11 中的 MQTT 和 3.12 中的原生 MQTT 之间的延迟,我们使用 mqtt-bm-latency。
以下延迟基准测试是在一台具有 8 个 CPU 和 32 GB RAM 的物理 Ubuntu 22.04 机器上执行的。客户端和单节点 RabbitMQ 服务器运行在同一台机器上。
在 rabbitmq-server 的根目录中,使用 4 个调度程序线程启动服务器:
make run-broker PLUGINS="rabbitmq_mqtt" RABBITMQ_CONFIG_FILE="rabbitmq.conf" RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS="+S 4"
服务器使用以下 rabbitmq.conf 运行:
mqtt.mailbox_soft_limit = 0
mqtt.tcp_listen_options.nodelay = true
mqtt.tcp_listen_options.backlog = 128
mqtt.tcp_listen_options.sndbuf = 87380
mqtt.tcp_listen_options.recbuf = 87380
mqtt.tcp_listen_options.buffer = 87380
classic_queue.default_version = 2
第一行仅在 RabbitMQ 3.12 中使用,它禁用了过载保护,确保所有 QoS 0 消息都传递给订阅者。
3.11 基准测试使用 Git 标签 v3.11.10。3.12 基准测试使用 Git 标签 v3.12.0-beta.2。
延迟和吞吐量 QoS 1
第一个基准测试比较发送给 QoS 1 订阅者的 QoS 1 消息的延迟和吞吐量。
./mqtt-bm-latency -clients 100 -count 10000 -pubqos 1 -subqos 1 -size 100 -keepalive 120 -topic t
该基准测试使用 100 个 MQTT 客户端“对”以 1:1 的拓扑结构进行发送。换句话说,共有 200 个 MQTT 客户端:100 个发布者和 100 个订阅者。每个发布者向单个订阅者发送 10,000 条消息。
所有 MQTT 客户端并发运行。但是,应用程序消息是从发布客户端同步发送到 RabbitMQ 的:每个发布者发送一条 QoS 1 消息,并等待直到收到来自 RabbitMQ 的 PUBACK 后才发送下一条。
订阅客户端也会为收到的每个 PUBLISH 数据包回复一个 PUBACK 数据包。
上述命令中的主题 t 只是一个主题前缀。每个发布客户端通过在其索引附加到给定主题来发布到不同的主题(例如,第一个发布者发布到主题 t-0,第二个到 t-1,依此类推)。
3.11 的结果:
================= TOTAL PUBLISHER (100) =================
Total Publish Success Ratio: 100.000% (1000000/1000000)
Total Runtime (sec): 75.835
Average Runtime (sec): 75.488
Pub time min (ms): 0.167
Pub time max (ms): 101.331
Pub time mean mean (ms): 7.532
Pub time mean std (ms): 0.037
Average Bandwidth (msg/sec): 132.475
Total Bandwidth (msg/sec): 13247.470
================= TOTAL SUBSCRIBER (100) =================
Total Forward Success Ratio: 100.000% (1000000/1000000)
Forward latency min (ms): 0.177
Forward latency max (ms): 94.045
Forward latency mean std (ms): 0.032
Total Mean forward latency (ms): 6.737
3.12 的结果:
================= TOTAL PUBLISHER (100) =================
Total Publish Success Ratio: 100.000% (1000000/1000000)
Total Runtime (sec): 55.955
Average Runtime (sec): 55.570
Pub time min (ms): 0.175
Pub time max (ms): 96.725
Pub time mean mean (ms): 5.550
Pub time mean std (ms): 0.031
Average Bandwidth (msg/sec): 179.959
Total Bandwidth (msg/sec): 17995.888
================= TOTAL SUBSCRIBER (100) =================
Total Forward Success Ratio: 100.000% (1000000/1000000)
Forward latency min (ms): 0.108
Forward latency max (ms): 51.421
Forward latency mean std (ms): 0.028
Total Mean forward latency (ms): 2.880
Total Publish Success Ratio 和 Total Forward Success Ratio 验证了 RabbitMQ 总共处理了 100 个发布者 * 10,000 条消息 = 1,000,000 条消息。
Total Runtime 和 Average Runtime 显示了第一个巨大的性能差异:使用原生 MQTT,每个发布者在 55 秒内从服务器接收到所有 10,000 条发布者确认 (PUBACK)。3.11 中的 MQTT 需要 75 秒。
Pub time mean mean 表明,平均发布者确认延迟从 3.11 中的 7.532 毫秒降低了 1.982 毫秒或 26%,降至 3.12 中的 5.550 毫秒。
如图 7 所示,100 个并发客户端每个同步向 RabbitMQ 发布 MQTT QoS 1 消息的吞吐量从 3.11 中的每秒 13,247 条消息增加了 4,748 条消息/秒或 36%,达到 3.12 中的每秒 17,995 条消息。
forward latency(转发延迟)统计信息代表端到端延迟,即从发布消息到客户端接收到消息需要多长时间。它是通过让发布者在每个消息有效负载中包含时间戳来测量的。
图 8 说明 Total mean forward latency(总平均转发延迟)从 3.11 中的 6.737 毫秒降低了 3.857 毫秒或 57%,降至 3.12 中的 2.880 毫秒。
延迟 QoS 0
第二个延迟基准测试看起来与第一个非常相似,但对发布的消息使用 QoS 0,并对订阅使用 QoS 0。
./mqtt-bm-latency -clients 100 -count 10000 -pubqos 0 -subqos 0 -size 100 -keepalive 120 -topic t
由于客户端以 CleanSession=1 连接,RabbitMQ 会为每个订阅者创建一个 rabbit_mqtt_qos0_queue。因此跳过了队列进程。
此处省略了 3.11 和 3.12 的发布者结果,因为它们非常低(Pub time mean mean 为 5 微秒)。发布客户端立即将所有 100 万条消息发送(“即发即弃”)到服务器,而无需等待 PUBACK 响应。
3.11 的结果:
================= TOTAL SUBSCRIBER (100) =================
Total Forward Success Ratio: 100.000% (1000000/1000000)
Forward latency min (ms): 3.907
Forward latency max (ms): 12855.600
Forward latency mean std (ms): 883.343
Total Mean forward latency (ms): 7260.071
3.12 的结果:
================= TOTAL SUBSCRIBER (100) =================
Total Forward Success Ratio: 100.000% (1000000/1000000)
Forward latency min (ms): 0.461
Forward latency max (ms): 5003.936
Forward latency mean std (ms): 596.133
Total Mean forward latency (ms): 2426.867
Total mean forward latency 从 3.11 中的 7,260.071 毫秒降低了 4,833.204 毫秒或 66%,降至 3.12 中的 2,426.867 毫秒。
与 QoS 1 基准测试相比,3.11 中 7.2 秒和 3.12 中 2.4 秒的端到端延迟非常高,因为代理同时被所有 100 万条 MQTT 消息暂时淹没。
在 advanced.config 中禁用信用流控(通过设置 {credit_flow_default_credit, {0, 0}})并不能改善 3.11 中的延迟。
感兴趣的读者可以尝试对 3.12 运行相同的基准测试,同时禁用功能标志 rabbit_mqtt_qos0_queue,以查看跳过队列进程带来的延迟差异。
3.12 中的原生 MQTT 还有什么改进?
以下列表包含了 3.12 中原生 MQTT 的其他变化项:
- 实现了 MQTT 3.1.1 允许 SUBACK 数据包包含失败返回代码 (0x80) 的功能。
- 删除了 AMQP 0.9.1 头
x-mqtt-dup,因为其在 3.11 及之前的 MQTT 插件中的用法是错误的。如果您的 AMQP 0.9.1 客户端依赖此头,这是一个重大变更。 - MQTT 插件创建的所有队列都是持久化的。这样做是为了方便在未来的 RabbitMQ 版本中过渡到 Khepri。
- MQTT 插件为干净会话创建排他队列(而不是自动删除队列)。
- 在 pg 中实现了 MQTT 客户端 ID 跟踪(以便在新客户端以相同客户端 ID 连接时终止现有连接)。当启用功能标志
delete_ra_cluster_mqtt_node时,先前用于跟踪 MQTT 客户端 ID 的 Ra 集群将被删除。旧的 Ra 集群需要大量内存,并且在同时断开过多客户端时成为瓶颈。 - 使用 全局计数器 实现了 Prometheus 指标。协议标签的值为
mqtt310或mqtt311。此类指标的一个示例是:rabbitmq_global_messages_routed_total{protocol="mqtt311"} 10。 - 优化了 MQTT 解析器。
- 从未发布应用程序消息的 MQTT 客户端在内存和磁盘告警期间不会被阻塞,因此订阅者可以继续清空队列。
我应该将 RabbitMQ 用作 MQTT 代理吗?
原生 MQTT 开启了许多新的用例,因为它允许数量级更多的物联网设备连接到 RabbitMQ。用例范围广泛。
RabbitMQ 只是市场上的一个 MQTT 代理。还有其他优秀的 MQTT 代理,其中一些能够处理比 RabbitMQ 更多的 MQTT 客户端连接,因为其他代理专门针对 MQTT 进行优化。
RabbitMQ 的优势在于它不仅仅是一个 MQTT 代理。它是一个通用的多协议代理。RabbitMQ 的真正力量在于协议互操作性,结合灵活的路由和队列类型选择。
举个例子:作为一家支付处理公司,您可以将世界各地分布的数十万或几百万台现金终端连接到中央 RabbitMQ 集群,每台终端每隔几分钟向 RabbitMQ 发送一条 MQTT QoS 1 消息(每当客户进行支付时)。
这些包含支付数据的 MQTT 消息对您的业务至关重要。在任何情况下(节点故障、磁盘故障、网络分区),消息都不允许丢失。它们还需要由一个或多个微服务处理。
一种解决方案是将几个仲裁队列绑定到主题交换机,每个仲裁队列绑定不同的路由键,从而将负载分散到各队列。每个仲裁队列通过 Raft 一致性算法在 3 或 5 个 RabbitMQ 节点之间复制,提供高数据安全性。每个仲裁队列中的消息可以由不同的微服务处理,这些微服务可以使用比 MQTT 更专业的通信协议,例如 AMQP 0.9.1 或 AMQP 1.0。
或者,MQTT 插件可以将 MQTT 支付消息转发到在多个 RabbitMQ 节点上复制的 超级流 (super stream)。然后,其他客户端可以通过使用 RabbitMQ Streams 协议多次从流中消费相同的消息来进行不同形式的数据分析。
RabbitMQ 提供的灵活性几乎是无限的,市场上没有其他消息代理能提供如此广泛的协议互操作性、数据安全性、容错性、可扩展性和性能。
未来工作
3.12 中发布的原生 MQTT 是 RabbitMQ 用作 MQTT 代理的关键一步。然而,MQTT 之旅并没有在这里结束。将来,我们计划(不作承诺!)添加以下 MQTT 改进:
- 添加对 MQTT 5.0 (v5) 的支持——这是社区强烈要求的功能。您可以在 PR #7263 中跟踪当前进度。截至 3.12 版本,RabbitMQ 仅支持 MQTT 3.1.1 (v4) 和 MQTT 3.1.0 (v3)。原生 MQTT 是开发 MQTT 5.0 支持的先决条件。
- 在升级处理数百万个 MQTT 连接的 RabbitMQ 节点时提高弹性,并提高大规模客户端断开连接(例如由于负载均衡器崩溃)时的弹性。
- 支持大规模扇出 QoS 1,将消息至少一次发送给所有设备。新的
rabbit_mqtt_qos0_queue类型通过不在发送和接收 Erlang 进程中保持任何状态来改善大规模扇出。然而,截至今日,向数百万个队列发送需要队列确认的消息(MQTT 中的 QoS 1 或 AMQP 0.9.1 中的发布者确认)应极其谨慎:始终使用相同的发布者,并且操作要非常罕见(每几分钟一次)。请注意,1:1 拓扑和大规模扇入对 RabbitMQ 来说并不成问题。
总结
尽管 RabbitMQ 在 3.11 中支持 MQTT,但在客户端连接数量方面的可扩展性很差,内存使用量过大。
3.12 中发布的原生 MQTT 将 RabbitMQ 转变为一个 MQTT 代理。它允许将数百万个客户端连接到 RabbitMQ。即使您不打算连接那么多客户端,通过将您的 MQTT 工作负载升级到 3.12,您也将大幅节省基础设施成本,因为内存使用量降低了高达 95%。
作为下一步,我们计划添加对 MQTT 5.0 的支持。
