You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在MQTT事件监听器中保持事务开启以避免LazyInitializationException

解决MQTT回调方法中Spring事务不生效的问题

你当前的问题是messageArrived()方法上的@Transactional注解不生效,导致操作plant.getSoilHumidityList()时触发LazyInitializationException——这是因为MQTT客户端回调的方法直接调用原始类实例,没有经过Spring的事务代理。以下是两种可行的解决方案:

方案1:使用TransactionTemplate手动管理事务

直接在回调方法中通过TransactionTemplate显式控制事务,绕过代理问题:

@Service
@Slf4j
@RequiredArgsConstructor
public class SoilHumidityEventListener {

  private final MqttClient mqttClient;
  private final ObjectMapper objectMapper;
  private final PlantService plantService;
  private final TransactionTemplate transactionTemplate; // 注入事务模板

  @Value("${mqtt.soil_humidity_topic}")
  private String soilHumidityTopic;

  // 订阅操作无需事务,移除@Transactional
  public void subscribe() throws MqttException {
    mqttClient.subscribe(soilHumidityTopic, this::messageArrived);
  }

  public void messageArrived(String topic, MqttMessage message) {
    try {
      transactionTemplate.execute(status -> {
        SoilHumidityDto soilHumidityDto =
            objectMapper.readValue(message.toString(), SoilHumidityDto.class);
        Plant plant = plantService.findById(soilHumidityDto.getPlantId());
        SoilHumidity soilHumidity = SoilHumidity.fromDto(soilHumidityDto);
        soilHumidity.setPlant(plant);
        plant.getSoilHumidityList().add(soilHumidity);
        plantService.save(plant);
        return null;
      });
    } catch (Exception e) {
      log.error("Error parsing and saving soil humidity message: {}", e.getMessage());
    }
  }
}

方案2:抽离业务逻辑到独立的事务方法

将核心业务逻辑抽成单独的@Transactional方法,通过Spring代理对象调用该方法:

@Service
@Slf4j
@RequiredArgsConstructor
public class SoilHumidityEventListener {

  private final MqttClient mqttClient;
  private final ObjectMapper objectMapper;
  private final PlantService plantService;
  private final SoilHumidityEventListener self; // 注入自身代理对象

  @Value("${mqtt.soil_humidity_topic}")
  private String soilHumidityTopic;

  public void subscribe() throws MqttException {
    mqttClient.subscribe(soilHumidityTopic, this::messageArrived);
  }

  public void messageArrived(String topic, MqttMessage message) {
    try {
      SoilHumidityDto soilHumidityDto =
          objectMapper.readValue(message.toString(), SoilHumidityDto.class);
      self.processSoilHumidity(soilHumidityDto); // 通过代理触发事务
    } catch (Exception e) {
      log.error("Error parsing and saving soil humidity message: {}", e.getMessage());
    }
  }

  @Transactional
  public void processSoilHumidity(SoilHumidityDto soilHumidityDto) {
    Plant plant = plantService.findById(soilHumidityDto.getPlantId());
    SoilHumidity soilHumidity = SoilHumidity.fromDto(soilHumidityDto);
    soilHumidity.setPlant(plant);
    plant.getSoilHumidityList().add(soilHumidity);
    plantService.save(plant);
  }
}

问题根源说明

原来的@Transactional注解失效是因为:

  • 订阅MQTT主题时,你传递的是this::messageArrived,这里的this是SoilHumidityEventListener的原始实例,而非Spring生成的事务代理实例。
  • Spring的事务是通过AOP代理实现的,只有通过代理对象调用方法时,才会触发事务拦截器。回调方法直接调用原始实例方法,自然不会开启事务。

内容的提问来源于stack exchange,提问作者MehdiB

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 14:22:46