mqtt-plus 架构解析(二):一条 MQTT 消息如何到达你的 @MqttListener

摘要

很多人第一次使用 @MqttListener 时,感受到的是“写起来很顺”,但框架内部到底做了哪些事,往往并不直观。本文沿着 MqttInboundMessageSink -> MqttMessageRouter -> MqttListenerRegistry -> ListenerInvoker 这条链路,拆解一条 MQTT 消息是如何到达监听方法的,并重点解释两个关键设计:为什么反序列化发生在 listener 分发之后,以及为什么用户代码不能直接跑在客户端回调线程里。


项目地址

项目地址:

https://github.com/mqttplus/mqtt-plus

配套的示例工程:

https://github.com/mqttplus/mqtt-plus-examples

如果你对这个方向感兴趣,欢迎关注、试用,也欢迎一起交流 issue 和 PR。

如果这篇文章对你有帮助,欢迎点赞、收藏,也欢迎给项目一个 Star。


假设你写了这样一段代码:

@MqttListener(broker = "primary", topics = "drone/+/status", payloadType = DroneStatus.class)
public void onStatus(DroneStatus status, MqttHeaders headers, @MqttTopic String topic) {
    // business logic
}

消息来了,topic 是 drone/001/status,payload 是一段 bytes。对使用者来说,这件事看起来像“框架自动把消息送到了方法里”;但从框架设计角度看,这个过程至少要回答四个问题:

  • 消息先交给谁?
  • topic 是怎么匹配到 listener 的?
  • payload 是在什么时机被转换的?
  • 方法调用为什么不能直接跑在 MQTT 客户端回调线程里?

这一篇文章就是把这四个问题拆开讲清楚。


一、这篇文章到底想回答什么?

这一篇只回答三个问题:

  • 一条消息从 broker 进入框架后,会经过哪些核心组件
  • 为什么一条消息可能匹配多个 listener,而且每个 listener 都要独立转换和独立调用
  • 为什么 mqtt-plus 选择“先路由、再转换、再调用”,而不是提前把 payload 统一反序列化好

如果只先记住一句话,可以记这句:

mqtt-plus 的路由链不是“收到消息就直接调方法”,而是先用 brokerId + topic 找到 listener,再按 listener 的 payloadType 独立完成转换和调用。

这个顺序看起来多了一层,但它恰好决定了多 listener、错误处理和扩展点是否还能成立。


二、先看整条消息流转链路

先把主链路拉平。只要这张图看明白了,后面的 topic 匹配、payload 转换和线程模型都会顺很多。

Listener ListenerInvoker PayloadConverter Interceptor MqttListenerRegistry MqttMessageRouter MqttInboundMessageSink Adapter Broker Listener ListenerInvoker PayloadConverter Interceptor MqttListenerRegistry MqttMessageRouter MqttInboundMessageSink Adapter Broker deliver message onMessage(brokerId, topic, payload, headers) route(...) resolve(brokerId, topic) matched listeners beforeHandle() convert(payload, targetType) invoke(definition, convertedPayload, context) call @MqttListener method afterHandle()

在源码里,这条链的几个关键锚点很明确:

  • MqttInboundMessageSink 定义了统一入站入口:onMessage(String brokerId, String topic, byte[] payload, MqttHeaders headers)
  • MqttMessageRouter 只暴露一个核心动作:route(...)
  • DefaultMqttMessageRouter 负责把“查 listener、执行 interceptor、转换 payload、调用方法、聚合错误动作”这些步骤串起来
  • MqttListenerRegistry 负责基于 brokerId + topic 找匹配的 listener
  • ListenerInvoker 负责把最终参数真正送进方法

也就是说,mqtt-plus 不是把“监听方法调用”塞进 adapter 里完成的,而是先把 adapter 收到的消息提升成一个框架内部统一的入站模型,然后交给路由器处理。这一步很关键,因为只有这样,多个 adapter 才能共享同一套路由逻辑。


三、topic 是怎么匹配到 listener 的?

topic 匹配发生在 MqttListenerRegistry.resolve(brokerId, topic) 这一层。

从实现上看,MqttListenerRegistry 内部维护的是 CopyOnWriteArrayList<MqttListenerDefinition>。每个 MqttListenerDefinition 里都带着 broker、topics、qos、payloadType、bean、method 等元信息。路由时,它会先做两层过滤:

  1. 先看 broker 是否匹配。只有 definition.getBroker().equals(brokerId) 或者监听器写的是 *,才继续往下。
  2. 再遍历这个 listener 声明的 topic pattern,交给 MqttTopicMatcher.matches(pattern, topic) 去判断。

MqttTopicMatcher 的逻辑不复杂,但它把几个关键语义都定死了:

  • + 只匹配单层 topic
  • # 只在最后一层时表示“后续全部匹配”
  • 如果 topic 以 $ 开头,而订阅模式不以 $ 开头,则直接不匹配

可以用一张简单的对照图快速建立直觉:

订阅模式消息 topic是否匹配原因
drone/+/statusdrone/001/status匹配+ 匹配单层
drone/+/statusdrone/001/telemetry不匹配最后一层不同
alert/#alert/warn/device/battery匹配# 匹配后续所有层
alert/#device/alert/warn不匹配起始层不同

这里还有一个容易忽略但对后面很重要的点:resolve 的结果不是单个 listener,而是一个列表。

这意味着,mqtt-plus 从一开始就把“同一条消息匹配多个 listener”视为正常情况,而不是异常情况。也正因为如此,后面才会有独立转换、独立调用和错误聚合这些设计。


四、一条消息,为什么可能会匹配多个 listener?

如果只从业务角度看,很多人会下意识觉得“一条消息应该只交给一个处理器”。但在注解驱动模型下,这个假设并不成立。

比如同一条消息 drone/001/status,完全可能同时被下面几类 listener 消费:

  • 一个用 String 接收原始文本做日志记录
  • 一个用 DroneStatus 接收结构化对象做业务处理
  • 一个用 byte[] 接收原始 payload 做审计或转发

这也是为什么 DefaultMqttMessageRouter.route(...) 的第一步不是“先转换 payload”,而是先拿到 matches,然后逐个 listener 去处理。

Inbound MQTT Message
topic=drone/001/status

MqttListenerRegistry.resolve()

Listener A
payloadType=String

Listener B
payloadType=DroneStatus

Listener C
payloadType=byte[]

String converter

JSON converter

Raw bytes

invoke A

invoke B

invoke C

这张图想表达的核心不是“框架多做了几步”,而是:同一条消息在进入不同 listener 时,其实可能对应完全不同的参数形态。

也正因为如此,mqtt-plus 的路由器在循环每个 MqttListenerDefinition 时,都会重新创建 MqttContext、重新选择 PayloadConverter、重新走一次 ListenerInvoker.invoke(...)。这是独立分发模型的代价,但也是它能支撑灵活监听签名的原因。

设计决策: 反序列化不是在消息进入框架的那一刻统一做掉,而是在 listener 分发之后按 payloadType 独立执行。因为在 mqtt-plus 的模型里,同一条消息本来就可能被不同 payloadType 的 listener 同时消费。


五、为什么真正的方法调用不是直接 method.invoke(payload)

从 core 的角度看,ListenerInvoker 只是一个抽象:给我一个 MqttListenerDefinition、一个已转换的 payload 和一个 MqttContext,我负责把它变成真正的方法调用。

这个设计本身就说明了一件事:路由器并不关心方法参数怎么拼出来,它只关心“路由到谁”。

在非 Spring 场景里,测试里有一个很简单的 ReflectiveListenerInvoker,逻辑基本就是:

  • 无参方法就直接 invoke(bean)
  • 单参方法就把 payload 塞进去

但在 Spring 环境里,starter 默认装配的是 SpringMqttListenerInvoker。它会把 converted payload + raw payload + topic + headers 一起交给 MqttListenerMethodArgumentResolver,然后再决定每个参数位该填什么。

这一步正好解释了为什么在 @MqttListener 方法里可以同时拿到:

  • 业务 payload
  • MqttHeaders
  • @MqttTopic 标注的 topic
  • 原始 byte[]

也就是说,ListenerInvoker 的存在不是为了“包一层反射显得高级”,而是为了把“路由”和“方法参数绑定”这两件事拆开。

这条边界一旦立住,路由器就不需要知道 Spring 方法参数解析细节,而 Spring 侧也不需要碰消息匹配逻辑。


六、为什么用户代码不能直接跑在客户端回调线程里?

这个问题经常被低估,但它其实是整个消息链里最工程化的一层。

PahoMqttClientAdapter 里,真正的消息入口不是直接调用 inboundMessageSink.onMessage(...),而是先把消息交给 inboundExecutor.submit(...)

  • messageArrived(...)
  • handleMessage(topic, message)
  • inboundExecutor.submit(() -> inboundMessageSink.onMessage(...))

这个设计非常直接地表达了一个原则:用户代码不能跑在底层 MQTT 客户端的回调线程里。

原因有三个:

  1. 如果 listener 处理慢,直接阻塞客户端回调线程,会影响后续消息接收。
  2. 如果 listener 抛异常,把客户端线程拖进不可控状态,故障边界会变得很模糊。
  3. 一旦要做线程池配置、削峰、隔离和观测,单纯依赖客户端回调线程几乎没有回旋空间。

starter 里之所以还要把 inboundThreadPool 放进 MqttBrokerDefinition 和配置属性里,本质上也是在承认:入站线程池是 broker 级别的工程配置,而不是某个 adapter 的私有细节。

设计决策: 用户 listener 绝不能直接运行在 MQTT 客户端回调线程里。mqtt-plus 选择先把消息移交给入站执行器,再进入统一路由链,这不是“多绕一步”,而是把 IO 层和业务层真正隔离开。

这里顺手也能理解一个对比:PahoMqttClientAdapter 把“线程切换”写得很显式,而 SpringIntegrationMqttClientAdapter 把更多调度能力交给 Spring Integration 基础设施。但不管外表怎么不同,进入 MqttInboundMessageSink 之后,后面的路由链仍然保持统一。


七、小结

这一篇真正想讲清的是:mqtt-plus 的消息路由不是一个“收到消息后立刻调方法”的黑盒,而是一条明确的、可扩展的处理链。

如果把结论压缩一下,可以记住这几件事:

  • MqttInboundMessageSink 是统一入站入口,adapter 负责把底层客户端消息送到这里。
  • DefaultMqttMessageRouter 真正串起了 resolve -> interceptor -> convert -> invoke -> aggregate 这条主链。
  • MqttListenerRegistryMqttTopicMatcher 共同决定“这条消息该交给哪些 listener”。
  • ListenerInvoker 把“路由到谁”和“方法参数怎么绑定”拆开,给 Spring 风格的监听方法留下了空间。
  • “先路由、再转换”与“不要跑在客户端回调线程里”这两个决定,看起来像实现细节,其实是后续错误处理、拦截器和多 adapter 统一行为的前提。

下一篇会继续沿着这条主线往前走,但聚焦的问题会更窄一些:为什么 mqtt-plus 要把序列化和反序列化拆成两条独立的链,以及这种拆分到底换来了什么。


系列导航

本文是 mqtt-plus 架构解析 系列的第 2/10 篇。

#主题链接
1总览:分层架构与设计哲学链接
2消息路由:一条 MQTT 消息如何到达你的 @MqttListener本文
3Payload 序列化与反序列化:双链设计的取舍链接
4拦截器链:MqttMessageInterceptor 的扩展点设计链接
5错误处理:ErrorAction 聚合策略的设计逻辑链接
6多 Broker 管理:如何让一个应用同时连接多个 MQTT 服务链接
7动态订阅与重连恢复:Reconciler 的协调机制链接
8Spring Boot 自动装配:零件是怎么被粘合起来的链接
9测试体系:MqttTestTemplateEmbeddedBroker 的设计链接
10从内部项目到开源框架:mqtt-plus 的抽取过程与决策链接

上一篇:总览:分层架构与设计哲学
下一篇:Payload 序列化与反序列化:双链设计的取舍

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐