跳转至主要内容
版本: 5.0

RocketMQ MQTT 概述

传统的消息队列 MQ 主要用于服务(端)与服务之间的消息通信,例如电商领域的交易消息、支付消息、物流消息等。然而,在消息的大类之下,还有另一个非常重要且常见的消息领域,即物联网(IoT)终端设备消息。近年来,我们看到了智能家居和工业互联带来的面向 IoT 设备的消息爆发式增长,而发展了十多年的移动互联网 APP 端消息量级依然巨大。终端设备的消息量级比传统服务器大几个数量级,且仍在快速增长。

如果有一个统一的消息系统(产品)来提供多场景计算(如流、事件)和多场景(IoT、APP)接入,实际上是非常有价值的,因为消息本身也是重要的数据。通过单一系统,可以最小化存储成本,并有效避免因不同系统间数据同步而导致的各种一致性问题和挑战。

image

基于此,我们推出了 RocketMQ-MQTT 扩展项目,旨在实现 RocketMQ 对 IoT 设备和服务器消息的统一接入,并提供集成的消息存储和互通能力。

MQTT 协议介绍

在 IoT 终端场景中,MQTT 协议被行业广泛采用。MQTT 起源于 IoT 环境,是一种基于发布/订阅(Pub/Sub)模型的轻量级消息传输协议,专为低带宽和不可靠网络环境设计。它最初由 IBM 开发,现由 OASIS 联盟维护为开放标准。MQTT 被广泛应用于物联网、智能硬件、车联网、智慧城市、远程医疗、电力、石油及能源等领域。

其核心通信模型也是发布/订阅(Pub/Sub),与 RocketMQ 类似。但它在订阅模式上提供了更大的灵活性,支持多级主题订阅(如 /t/t1/t2)和通配符订阅(如 /t/t1/+)。使用 MQTT,可以轻松实现消息的广播、组播和单播。

RocketMQ MQTT 架构设计

RocketMQ MQTT 架构的目标是实现消息存储和分发的统一管理,在不侵入 RocketMQ Broker 核心逻辑的前提下实现多协议集成。为此,我们设计了两个基础模型:队列存储模型和推拉模型。

队列存储模型

image

我们设计了一种用于多维度分发的主题队列模型。如上图所示,消息可以来自各种接入场景(如服务器端的 MQ/AMQP 和客户端的 MQTT),但只会写入并存储一份到 CommitLog 中,然后分发到多个需求场景的队列索引(ConsumerQueue)。例如,服务器端场景(MQ/AMQP)可以根据一级 Topic 队列进行传统的服务器端消费,而客户端 MQTT 场景则可以根据 MQTT 的多级 Topic 和通配符订阅进行消费。

这样的队列模型可以同时支持服务器和终端场景的接入以及消息发送与接收,从而实现集成化的目标。

推拉模型

image

上图展示了推拉模型。图中的 P 节点是协议网关或 Broker 插件,终端设备通过 MQTT 协议连接到网关节点。消息可以从各种场景(MQ/AMQP/MQTT)发送。在存储到 Topic 队列后,会有一个通知(Notify)逻辑模块实时感知新消息的到达,然后生成一个消息事件(即消息的主题名称)。该事件被推送给网关节点,网关节点根据已连接终端设备的订阅状态进行内部匹配,找出哪些终端设备可以匹配,随后触发对存储层的拉取请求,读取消息并推送给终端设备。

架构概述

image 我们的目标是基于 RocketMQ 实现一个集成且自闭环的系统,但不希望 Broker 侵入过多的场景逻辑。我们抽象出一个协议计算层,它可以是网关或 Broker 插件。Broker 专注于解决队列问题,并进行一些队列存储的适配或转换,以满足上述计算需求。协议计算层负责协议接入,并且必须是可插拔和可部署的。