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

SpringBoot中无法将MQTT Broker主题消息传入ObjectMapper解析

解决方案与代码优化建议

1. 重构MQTT消息接收逻辑(替换现有MqttBeans)

放弃手动管理MQTT通道的方式,改用Spring官方的spring-integration-mqtt依赖,简化连接配置与消息流转,从根源解决payload传递问题。

第一步:添加依赖(pom.xml)

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-integration</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-mqtt</artifactId>
</dependency>

第二步:配置MQTT连接与消息处理通道

创建MqttConfig类,统一管理客户端连接、主题监听和消息处理器:

import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

@Configuration
public class MqttConfig {

    @Value("${mqtt.broker.url}")
    private String brokerUrl;

    @Value("${mqtt.client.id}")
    private String clientId;

    @Value("${mqtt.topic.weather-data}")
    private String weatherTopic;

    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        factory.setServerURIs(brokerUrl);
        // 如需认证,添加用户名密码配置
        // factory.setUserName("your-username");
        // factory.setPassword("your-password");
        return factory;
    }

    @Bean
    public MessageChannel mqttInputChannel() {
        return new DirectChannel();
    }

    @Bean
    public MessageProducer inbound() {
        MqttPahoMessageDrivenChannelAdapter adapter =
                new MqttPahoMessageDrivenChannelAdapter(clientId + "-inbound",
                        mqttClientFactory(), weatherTopic);
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1);
        adapter.setOutputChannel(mqttInputChannel());
        return adapter;
    }

    // 直接绑定消息处理器,自动接收并解析payload
    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel")
    public MessageHandler weatherMessageHandler(ObjectMapper objectMapper, WeatherRepository weatherRepository) {
        return message -> {
            String payload = (String) message.getPayload();
            try {
                Weather weather = objectMapper.readValue(payload, Weather.class);
                weatherRepository.save(weather);
            } catch (JsonProcessingException e) {
                // 可替换为日志记录或异常告警逻辑
                e.printStackTrace();
            }
        };
    }
}

2. 实体类与数据库层定义

确保Weather实体与Postgres表字段一一对应,用Spring Data JPA简化数据库操作:

import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import lombok.Data;

@Entity
@Data
public class Weather {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;
    private String temperature;
    private String humidity;
    private String windSpeed;
    // 根据实际JSON字段添加对应属性
}
import org.springframework.data.jpa.repository.JpaRepository;

public interface WeatherRepository extends JpaRepository<Weather, Long> {
}

3. 替换StartupUtility为单元测试

原来的硬编码测试逻辑改为单元测试,专门验证JSON解析能力:

import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;

@SpringBootTest
public class WeatherParsingTest {

    @Autowired
    private ObjectMapper objectMapper;

    @Test
    public void testJsonParsing() throws Exception {
        String testJson = "{\"temperature\":\"25°C\",\"humidity\":\"60%\",\"windSpeed\":\"10km/h\"}";
        Weather weather = objectMapper.readValue(testJson, Weather.class);
        assert weather.getTemperature().equals("25°C");
        assert weather.getHumidity().equals("60%");
    }
}

4. 常见问题排查

  • 检查MQTT客户端ID是否唯一,避免与其他服务冲突导致连接异常
  • 确认Mosquitto Broker开启了weather-data主题的订阅权限
  • 核对application.properties配置:
mqtt.broker.url=tcp://localhost:1883
mqtt.client.id=data-retriever-service
mqtt.topic.weather-data=weather-data

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 20:42:39