在开发充电桩订单服务时,我们遇到了一个值得深入思考的设计问题: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 MqttConsumersetCallbackconnect 的严格顺序,确保在 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(比如 OrderServiceOrderRepository 等),此时又该怎么办?

有几种常见的处理方式:

方案一: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 依赖 MqttConsumerMqttConsumer 依赖 MqttClient,Spring 对构造器循环依赖无解
时序控制 setCallback 必须严格在 connect 之前执行,Spring 的 Bean 初始化无法保证这种毫秒级的时序精度
职责单一 MqttConsumer 只是 MqttClient 的一个回调处理器,不需要被其他地方注入,进容器是画蛇添足

理解这个设计,不仅仅是知道"为什么不能这么做",更重要的是建立起对 Spring Bean 生命周期循环依赖异步回调时序这些核心概念的深层认知。这些知识在实际项目中会反复遇到,值得反复琢磨。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐