为什么 MqttConsumer 不交给 Spring 管理?——从一个充电桩项目说起
在开发充电桩订单服务时,我们遇到了一个值得深入思考的设计问题:
MqttConsumer作为 MQTT 消息的回调处理类,为什么没有加@Component注解交给 Spring 容器管理,而是选择手动new出来?这背后涉及到 Spring Bean 的生命周期、循环依赖、以及异步回调的时序控制等多个核心知识点。本文将从业务背景出发,逐层分析这个设计决策的原因。
一、背景:充电桩订单服务的通信架构
在充电桩项目中,订单服务需要和充电设备进行双向通信:
- 订单服务 → 设备:发送"开始充电"指令
- 设备 → 订单服务:返回"充电结果"(成功 or 失败)
这个通信是通过 EMQX(一个高性能的 MQTT 消息中间件)来实现的。整体流程如下:
订单服务 EMQX 充电设备
| | |
|--- 发布开始充电指令 ----->| |
| |<--- 订阅充电指令 Topic ----|
| |---- 推送开始充电指令 ----->|
| | |
|<--- 订阅充电结果 Topic ---| |
| |<--- 发布充电结果 ---------|
|<--- 推送充电结果 ---------| |
为了实现这套通信机制,订单服务中有两个关键组件:
MqttClient:MQTT 客户端,负责和 EMQX 建立连接、发布消息、订阅 Topic。MqttConsumer:MQTT 消息消费者(回调类),负责处理从 EMQX 收到的设备消息。
问题就出在这两个组件的关系上。
二、先搞清楚:MqttConsumer 是干什么的?
MqttConsumer 实现了 MqttCallbackExtended 接口,它本质上是一个事件回调处理器。
public class MqttConsumer implements MqttCallbackExtended {
private MqttClient mqttClient;
public MqttConsumer(MqttClient mqttClient) {
this.mqttClient = mqttClient;
}
/**
* 连接成功的回调
* 在这里订阅 Topic
*/
@Override
public void connectComplete(boolean reconnect, String serverURI) {
// 连接成功后,订阅设备返回充电结果的 Topic
mqttClient.subscribe(MqttConstants.TOPIC_CHARGING_RESULT);
}
/**
* 消息到达的回调
* 在这里处理设备发回的消息
*/
@Override
public void messageArrived(String topic, MqttMessage message) throws Exception {
// 解析设备消息,处理充电结果
ChargingResultDto dto = JsonUtils.fromJson(message.toString(), ChargingResultDto.class);
if (ChargingConstants.CHARGING_RESULT_SUCCESS.equals(dto.getResult())) {
// 处理充电成功逻辑
} else {
// 处理充电失败逻辑
}
}
}
这个类里有一个非常关键的依赖:它需要 MqttClient 来订阅 Topic。
而 MqttClient 又需要 MqttConsumer 来设置回调:
mqttClient.setCallback(mqttConsumer); // MqttClient 需要 MqttConsumer
这就构成了一个经典的互相依赖关系。
三、如果交给 Spring 管理,会发生什么?
我们先假设把 MqttConsumer 加上 @Component,看看 Spring 会怎么处理。
3.1 Spring 的 Bean 创建流程回顾
Spring 容器在启动时,会按照以下流程创建和初始化 Bean:
1. 扫描所有 @Component、@Bean 等注解
2. 解析 Bean 之间的依赖关系
3. 按照依赖顺序,依次创建 Bean
4. 注入依赖(@Autowired / 构造器注入)
5. 调用初始化方法(@PostConstruct 等)
3.2 循环依赖问题
如果 MqttConsumer 是一个 Spring Bean,并且通过构造器注入 MqttClient:
@Component
public class MqttConsumer implements MqttCallbackExtended {
private final MqttClient mqttClient;
// 构造器注入
@Autowired
public MqttConsumer(MqttClient mqttClient) {
this.mqttClient = mqttClient;
}
}
同时,MqttClient 的 @Bean 方法在创建时又需要用到 MqttConsumer:
@Bean
public MqttClient mqttClient(MqttConsumer mqttConsumer) throws MqttException {
MqttClient mqttClient = new MqttClient(host, clientId);
mqttClient.setCallback(mqttConsumer); // 需要 mqttConsumer
mqttClient.connect(options);
return mqttClient;
}
此时,Spring 的依赖解析过程如下:
Spring 尝试创建 MqttClient Bean
└── MqttClient Bean 依赖 MqttConsumer Bean
└── 尝试创建 MqttConsumer Bean
└── MqttConsumer Bean 构造器依赖 MqttClient Bean
└── MqttClient Bean 正在创建中,还没完成!
💥 循环依赖,Spring 抛出异常!
Spring 对于构造器注入的循环依赖是无法解决的,会直接报错:
BeanCurrentlyInCreationException: Error creating bean with name 'mqttClient':
Requested bean is currently in creation: Is there an unresolvable circular reference?
为什么构造器注入无法解决循环依赖?
Spring 解决循环依赖的方式是"三级缓存"机制,通过提前暴露一个"半成品"的 Bean 引用来打破循环。但这个机制只对字段注入(@Autowired 字段)和Setter 注入有效。
构造器注入要求在创建对象时就必须传入所有依赖,如果依赖还没创建完,根本无法调用构造器,所以 Spring 对构造器注入的循环依赖束手无策。
3.3 即使用字段注入,也逃不过时序问题
也许有人会想,那我用字段注入(@Autowired 注解在字段上)来替代构造器注入,绕过循环依赖,可不可以?
@Component
public class MqttConsumer implements MqttCallbackExtended {
@Autowired // 字段注入,Spring 可以先创建半成品 Bean,再注入
private MqttClient mqttClient;
}
理论上 Spring 可以通过三级缓存机制解决这个循环依赖,两个 Bean 都能创建出来。但这里有一个更致命的问题:时序问题。
我们看看 MqttConfiguration 中 @Bean 方法里这几行代码的顺序:
@Bean
public MqttClient mqttClient() throws MqttException {
MqttClient mqttClient = new MqttClient(host, clientId);
// 第 1 步:必须先设置回调
mqttClient.setCallback(mqttConsumer);
// 第 2 步:再连接 EMQX
mqttClient.connect(options);
// ↑ connect() 方法执行后,连接成功的瞬间会触发 connectComplete() 回调!
// ↑ 如果此时 callback 没有设置好,connectComplete() 里的订阅 Topic 逻辑就会丢失!
return mqttClient;
}
mqttClient.connect(options) 是异步触发回调的:一旦和 EMQX 连接成功,会立刻调用 MqttConsumer.connectComplete() 方法。
而 connectComplete() 里的核心逻辑是订阅 Topic:
@Override
public void connectComplete(boolean reconnect, String serverURI) {
// 连接成功后订阅 Topic,如果此时 callback 没有设置,这段代码根本不会被执行!
mqttClient.subscribe(MqttConstants.TOPIC_CHARGING_RESULT);
}
如果 setCallback() 还没执行,connect() 就已经触发了连接成功事件,那么这次连接成功的 connectComplete 回调将永远丢失,Topic 就再也订阅不上了。
而如果 MqttConsumer 是 Spring Bean,Spring 什么时候完成字段注入、什么时候把 mqttConsumer 注入进来,这个时序是由 Spring 容器内部控制的,我们无法精确干预。
四、当前的正确做法:手动 new,掌控一切
理解了上面的两个问题(循环依赖 + 时序问题)之后,再来看当前的实现,就会觉得非常优雅:
@Bean
public MqttClient mqttClient() throws MqttException {
// 第 1 步:创建 MqttClient
MqttClient mqttClient = new MqttClient(host, clientId);
MqttConnectOptions options = new MqttConnectOptions();
options.setUserName(username);
options.setPassword(password.toCharArray());
// 第 2 步:手动 new MqttConsumer,把 mqttClient 通过构造器传入
// 此时 mqttConsumer 已经持有了 mqttClient 的引用
MqttConsumer mqttConsumer = new MqttConsumer(mqttClient);
// 第 3 步:把 mqttConsumer 设置为 mqttClient 的回调
// 此时 callback 已就位,随时可以响应连接成功事件
mqttClient.setCallback(mqttConsumer);
// 第 4 步:连接 EMQX
// 连接成功后,connectComplete() 会被立刻回调
// 此时 callback 已经设置好,connectComplete() 里的订阅逻辑可以正常执行
mqttClient.connect(options);
return mqttClient;
}
这套做法完美解决了两个问题:
解决循环依赖:MqttConsumer 不是 Spring Bean,它只是在 @Bean 方法里被手动 new 出来。MqttClient 的 @Bean 方法是一个普通的 Java 方法,内部想先创建什么、后创建什么,完全由我们自己决定,Spring 的依赖解析机制根本不会介入。
解决时序问题:手动编写代码,显式控制 new MqttConsumer → setCallback → connect 的严格顺序,确保在 connect() 被调用之前,回调已经就位。
五、深入理解:异步回调的时序本质
这里值得多说几句关于"异步回调时序"的问题,它是理解这个设计的核心。
MQTT 的连接过程是这样的:
mqttClient.connect(options)
↓
内部发起 TCP 连接请求(异步)
↓
连接建立成功
↓
触发 MqttConsumer.connectComplete() 回调(可能在另一个线程中)
↓
connectComplete() 里执行 mqttClient.subscribe(topic)
注意:connect() 方法本身可能在连接成功之前就返回了(取决于 MQTT 客户端的实现),connectComplete() 是在另一个线程中异步被调用的。
这意味着:
- 如果你先调用
connect(),再调用setCallback(),那么在这两个调用之间的窗口期内,如果连接恰好成功,connectComplete()就会在一个没有 callback 的客户端上触发,或者压根不触发,导致 Topic 订阅逻辑永久丢失。 - 网络速度越快(比如本地测试环境),这个 bug 越容易复现;网络较慢时可能看起来没问题,但这是一个隐患。
这种问题在并发和异步场景下非常常见,被称为 TOCTOU(Time of Check to Time of Use) 竞态条件。正确的做法就是确保在触发异步操作之前,所有的准备工作都已经就位,这也是为什么一定要先 setCallback,再 connect 的原因。
六、扩展思考:什么情况下应该交给 Spring 管理?
看到这里,可能有人会问:难道 MqttConsumer 就永远不能交给 Spring 管理吗?
不是的。如果 MqttConsumer 里需要依赖其他的 Spring Bean(比如 OrderService、OrderRepository 等),此时又该怎么办?
有几种常见的处理方式:
方案一:ApplicationContext 注入(懒获取)
让 MqttConsumer 实现 ApplicationContextAware,在需要用 Bean 的时候通过 ApplicationContext 手动获取:
public class MqttConsumer implements MqttCallbackExtended, ApplicationContextAware {
private ApplicationContext applicationContext;
@Override
public void setApplicationContext(ApplicationContext ctx) {
this.applicationContext = ctx;
}
@Override
public void messageArrived(String topic, MqttMessage message) {
// 在需要的时候才从容器中获取 Bean,避免循环依赖
OrderService orderService = applicationContext.getBean(OrderService.class);
orderService.handleChargingResult(...);
}
}
这种方式的好处是 MqttConsumer 本身还是手动 new,不参与 Spring 的 Bean 生命周期,彻底避免了循环依赖问题。缺点是稍微绕了一层,代码可读性略差。
方案二:构造器传入所有依赖
在 MqttConfiguration 的 @Bean 方法里,把 MqttConsumer 需要的其他 Spring Bean 通过方法参数注入进来,再传给 MqttConsumer 的构造器:
@Bean
public MqttClient mqttClient(OrderService orderService) throws MqttException {
MqttClient mqttClient = new MqttClient(host, clientId);
// 把 orderService 也传入 MqttConsumer
MqttConsumer mqttConsumer = new MqttConsumer(mqttClient, orderService);
mqttClient.setCallback(mqttConsumer);
mqttClient.connect(options);
return mqttClient;
}
这种方式更加直接,依赖关系一目了然,是推荐的做法。
方案三:使用 Spring 的 @Lazy 注解
如果坚持要把 MqttConsumer 注册为 Spring Bean,可以在一侧使用 @Lazy 注解,让 Spring 延迟初始化,从而打破循环:
@Bean
public MqttClient mqttClient(@Lazy MqttConsumer mqttConsumer) throws MqttException {
// ...
}
但这种方式治标不治本,时序问题依然存在,需要额外处理,不推荐在这个场景下使用。
七、好的编程习惯总结
通过这个案例,我们可以总结出几点实用的编程经验:
1. 并非所有的类都需要交给 Spring 管理。 Spring 容器管理的是需要在多处共享、需要享受 Spring 生命周期特性(如依赖注入、AOP、事务等)的对象。像 MqttConsumer 这样的回调类,它的生命周期完全依附于 MqttClient,由 MqttClient 负责持有和调用,不需要被其他地方注入,没有必要进容器。
2. 谁创建,谁负责。 MqttConsumer 是在 MqttConfiguration 的 @Bean 方法里创建的,它的生命周期就交给 MqttClient Bean 来负责。这种"就近创建、就近管理"的原则能有效避免对象归属不清的问题。
3. 异步场景下,准备工作必须在触发之前完成。 在任何异步回调的场景中,都要确保所有的"准备动作"(如 setCallback)在触发异步操作(如 connect)之前完成,避免竞态条件。
4. 循环依赖往往是设计问题的信号。 当你发现两个类互相依赖时,应该思考是否有设计上的问题。通常的解法是:引入第三方来解耦,或者重新审视职责划分,而不是想办法让 Spring "绕过"这个循环。
八、完整流程图
最后,用一张流程图来梳理整个过程:
Spring 容器启动
|
▼
MqttConfiguration.mqttClient() 方法执行
|
▼
new MqttClient(host, clientId) ← 创建 MQTT 客户端
|
▼
new MqttConsumer(mqttClient) ← 手动创建回调类,传入 mqttClient
|
▼
mqttClient.setCallback(mqttConsumer) ← 设置回调(必须在 connect 之前!)
|
▼
mqttClient.connect(options) ← 连接 EMQX
|
▼(异步,连接成功后触发)
MqttConsumer.connectComplete() ← 连接成功回调
|
▼
mqttClient.subscribe(TOPIC_CHARGING_RESULT) ← 订阅设备充电结果 Topic
|
▼(等待设备消息)
MqttConsumer.messageArrived() ← 收到设备消息回调
|
▼
解析 JSON → ChargingResultDto
|
├─ result == "start_success" → 保存成功订单记录
└─ result != "start_success" → 保存失败订单记录
总结
MqttConsumer 不交给 Spring 管理,核心原因可以归结为以下三点:
| 原因 | 具体说明 |
|---|---|
| 循环依赖 | MqttClient 依赖 MqttConsumer,MqttConsumer 依赖 MqttClient,Spring 对构造器循环依赖无解 |
| 时序控制 | setCallback 必须严格在 connect 之前执行,Spring 的 Bean 初始化无法保证这种毫秒级的时序精度 |
| 职责单一 | MqttConsumer 只是 MqttClient 的一个回调处理器,不需要被其他地方注入,进容器是画蛇添足 |
理解这个设计,不仅仅是知道"为什么不能这么做",更重要的是建立起对 Spring Bean 生命周期、循环依赖、异步回调时序这些核心概念的深层认知。这些知识在实际项目中会反复遇到,值得反复琢磨。
更多推荐




所有评论(0)