mqtt-plus 架构解析(二):一条 MQTT 消息如何到达你的 @MqttListener
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 转换和线程模型都会顺很多。
在源码里,这条链的几个关键锚点很明确:
MqttInboundMessageSink定义了统一入站入口:onMessage(String brokerId, String topic, byte[] payload, MqttHeaders headers)MqttMessageRouter只暴露一个核心动作:route(...)DefaultMqttMessageRouter负责把“查 listener、执行 interceptor、转换 payload、调用方法、聚合错误动作”这些步骤串起来MqttListenerRegistry负责基于brokerId + topic找匹配的 listenerListenerInvoker负责把最终参数真正送进方法
也就是说,mqtt-plus 不是把“监听方法调用”塞进 adapter 里完成的,而是先把 adapter 收到的消息提升成一个框架内部统一的入站模型,然后交给路由器处理。这一步很关键,因为只有这样,多个 adapter 才能共享同一套路由逻辑。
三、topic 是怎么匹配到 listener 的?
topic 匹配发生在 MqttListenerRegistry.resolve(brokerId, topic) 这一层。
从实现上看,MqttListenerRegistry 内部维护的是 CopyOnWriteArrayList<MqttListenerDefinition>。每个 MqttListenerDefinition 里都带着 broker、topics、qos、payloadType、bean、method 等元信息。路由时,它会先做两层过滤:
- 先看 broker 是否匹配。只有
definition.getBroker().equals(brokerId)或者监听器写的是*,才继续往下。 - 再遍历这个 listener 声明的 topic pattern,交给
MqttTopicMatcher.matches(pattern, topic)去判断。
MqttTopicMatcher 的逻辑不复杂,但它把几个关键语义都定死了:
+只匹配单层 topic#只在最后一层时表示“后续全部匹配”- 如果 topic 以
$开头,而订阅模式不以$开头,则直接不匹配
可以用一张简单的对照图快速建立直觉:
| 订阅模式 | 消息 topic | 是否匹配 | 原因 |
|---|---|---|---|
drone/+/status | drone/001/status | 匹配 | + 匹配单层 |
drone/+/status | drone/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 去处理。
这张图想表达的核心不是“框架多做了几步”,而是:同一条消息在进入不同 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 客户端的回调线程里。
原因有三个:
- 如果 listener 处理慢,直接阻塞客户端回调线程,会影响后续消息接收。
- 如果 listener 抛异常,把客户端线程拖进不可控状态,故障边界会变得很模糊。
- 一旦要做线程池配置、削峰、隔离和观测,单纯依赖客户端回调线程几乎没有回旋空间。
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这条主链。MqttListenerRegistry和MqttTopicMatcher共同决定“这条消息该交给哪些 listener”。ListenerInvoker把“路由到谁”和“方法参数怎么绑定”拆开,给 Spring 风格的监听方法留下了空间。- “先路由、再转换”与“不要跑在客户端回调线程里”这两个决定,看起来像实现细节,其实是后续错误处理、拦截器和多 adapter 统一行为的前提。
下一篇会继续沿着这条主线往前走,但聚焦的问题会更窄一些:为什么 mqtt-plus 要把序列化和反序列化拆成两条独立的链,以及这种拆分到底换来了什么。
系列导航
本文是 mqtt-plus 架构解析 系列的第 2/10 篇。
| # | 主题 | 链接 |
|---|---|---|
| 1 | 总览:分层架构与设计哲学 | 链接 |
| 2 | 消息路由:一条 MQTT 消息如何到达你的 @MqttListener | 本文 |
| 3 | Payload 序列化与反序列化:双链设计的取舍 | 链接 |
| 4 | 拦截器链:MqttMessageInterceptor 的扩展点设计 | 链接 |
| 5 | 错误处理:ErrorAction 聚合策略的设计逻辑 | 链接 |
| 6 | 多 Broker 管理:如何让一个应用同时连接多个 MQTT 服务 | 链接 |
| 7 | 动态订阅与重连恢复:Reconciler 的协调机制 | 链接 |
| 8 | Spring Boot 自动装配:零件是怎么被粘合起来的 | 链接 |
| 9 | 测试体系:MqttTestTemplate 与 EmbeddedBroker 的设计 | 链接 |
| 10 | 从内部项目到开源框架:mqtt-plus 的抽取过程与决策 | 链接 |
上一篇:总览:分层架构与设计哲学
下一篇:Payload 序列化与反序列化:双链设计的取舍
更多推荐
所有评论(0)