SpringBoot整合MQTT:后端订阅设备数据、解析传感器报文
SpringBoot整合MQTT后端订阅设备数据、解析传感器报文作者黒漂技术佬前面几篇文章都是在讲 MQTT 协议本身和硬件端的事情。现在该轮到后端程序员出场了——设备把数据发到了 MQTT Broker咱们的后端服务怎么把它接住这篇就手把手带你用 SpringBoot 整合 MQTT实现订阅传感器数据、解析报文、存入库的全流程。一、方案选型用哪个 MQTT 客户端库Java 生态里搞 MQTT 主要有两个选项方案说明推荐度org.eclipse.paho.client.mqttv3Eclipse Paho 原生客户端功能完整但偏底层⭐⭐⭐spring-integration-mqttSpring 官方集成方案基于 Paho 封装与 Spring 生态无缝对接⭐⭐⭐⭐⭐强烈推荐spring-integration-mqtt。理由很简单你已经在用 SpringBoot 了何必再从底层写一堆连接管理、线程池、异常处理的代码Spring Integration 把 MQTT 客户端包装成了 Spring 的MessageChannel和MessageHandler和你写 Controller 一个味儿。二、引入依赖和配置2.1 pom.xml!-- MQTT 核心依赖 --dependencygroupIdorg.springframework.integration/groupIdartifactIdspring-integration-mqtt/artifactId/dependency!-- JSON 处理 --dependencygroupIdcom.fasterxml.jackson.core/groupIdartifactIdjackson-databind/artifactId/dependencySpringBoot 的spring-boot-starter-integration会自动拉取 Integration 核心所以不需要额外引入。2.2 application.ymlmqtt:broker-url:tcp://192.168.1.100:1883client-id:${spring.application.name}-${random.value}username:adminpassword:admin123# 订阅的Topic列表topics:-agriculture///sensor/## QoS级别qos:1# 超时和心跳配置completion-timeout:3000keep-alive-interval:60# 是否异步发送async:trueclient-id里加了随机值是为了支持多实例部署——两个相同 client-id 的连接会互相踢下线你肯定不想这样。三、MQTT 配置类下面是一个可直接用于生产的 MQTT 配置类40行左右ConfigurationIntegrationComponentScanpublicclassMqttConfig{Value(${mqtt.broker-url})privateStringbrokerUrl;Value(${mqtt.client-id})privateStringclientId;Value(${mqtt.username})privateStringusername;Value(${mqtt.password})privateStringpassword;Value(${mqtt.completion-timeout})privateintcompletionTimeout;Value(${mqtt.keep-alive-interval})privateintkeepAliveInterval;Value(#{${mqtt.topics}.split(,)})privateListStringtopics;Value(${mqtt.qos})privateintqos;// ① 连接配置BeanpublicMqttConnectOptionsmqttConnectOptions(){MqttConnectOptionsoptionsnewMqttConnectOptions();options.setServerURIs(newString[]{brokerUrl});options.setUserName(username);options.setPassword(password.toCharArray());options.setCleanSession(false);// 持久会话options.setAutomaticReconnect(true);// 自动重连options.setKeepAliveInterval(keepAliveInterval);options.setConnectionTimeout(10);returnoptions;}// ② 客户端工厂BeanpublicMqttPahoClientFactorymqttClientFactory(){DefaultMqttPahoClientFactoryfactorynewDefaultMqttPahoClientFactory();factory.setConnectionOptions(mqttConnectOptions());returnfactory;}// ③ 入站通道MQTT Broker → 应用BeanpublicMessageChannelmqttInputChannel(){returnnewDirectChannel();}// ④ 入站适配器订阅TopicBeanpublicMessageProducerinbound(){MqttPahoMessageDrivenChannelAdapteradapternewMqttPahoMessageDrivenChannelAdapter(clientId,mqttClientFactory(),topics.toArray(newString[0]));adapter.setCompletionTimeout(completionTimeout);adapter.setConverter(newDefaultPahoMessageConverter());adapter.setQos(qos);adapter.setOutputChannel(mqttInputChannel());returnadapter;}}来逐段解读一下①MqttConnectOptions就像你上网时的连接设置。setAutomaticReconnect(true)告诉 Paho「断了就自己连回来别烦我」。②MqttPahoClientFactory工厂模式负责生产 MQTT 客户端实例。Spring Integration 内部会用它来创建连接。③DirectChannelSpring Integration 的消息通道简单理解就是一个「管道」消息从这里流进来。④MqttPahoMessageDrivenChannelAdapter入站适配器负责订阅 Topic把收到的消息灌入mqttInputChannel。topics.toArray(new String[0])支持多 Topic 订阅比如同时订阅温湿度 Topic 和光照 Topic。四、消息接收处理器配置写好了现在接收消息ComponentpublicclassSensorDataHandler{ServiceActivator(inputChannelmqttInputChannel)publicvoidhandleMessage(Message?message){// 获取 TopicStringtopic(String)message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC);// 获取 Payload消息体Stringpayload(String)message.getPayload();System.out.printf([收到消息] Topic: %s%n,topic);System.out.printf([消息内容] %s%n,payload);// 解析JSONtry{SensorDatadataparseSensorData(topic,payload);processSensorData(data);}catch(Exceptione){System.err.println(消息解析失败: e.getMessage());}}/** * 根据Topic路由到不同的解析逻辑 */privateSensorDataparseSensorData(Stringtopic,Stringpayload)throwsException{ObjectMappermappernewObjectMapper();if(topic.contains(/temperature)){returnmapper.readValue(payload,TemperatureData.class);}elseif(topic.contains(/humidity)){returnmapper.readValue(payload,HumidityData.class);}elseif(topic.contains(/light)){returnmapper.readValue(payload,LightData.class);}thrownewIllegalArgumentException(Unknown topic: topic);}privatevoidprocessSensorData(SensorDatadata){// 1. 数据校验if(!data.isValid()){log.warn(数据异常已丢弃: {},data);return;}// 2. 业务处理存库、告警、转发等sensorDataService.save(data);// 3. 实时推送WebSocket通知前端大屏webSocketService.push(data);}}ServiceActivator(inputChannel mqttInputChannel)这行是核心。它告诉 Spring「MQTT 来的消息从mqttInputChannel这个管道流过来交给handleMessage方法处理」。Spring Integration 的Message?对象封装了消息的 Header元数据和 Payload消息体。从 Header 里可以拿到 Topic、QoS、是否 Retained 等信息。五、消息发布后端向设备发指令光收不发怎么行咱们还得给设备下发控制指令开风机、关水泵之类的ServicepublicclassMqttCommandService{AutowiredprivateMqttPahoClientFactorymqttClientFactory;Value(${mqtt.client-id}-outbound)privateStringoutboundClientId;/** * 发送控制指令到指定设备 */publicvoidsendCommand(StringdeviceId,Stringcommand,MapString,Objectparams){StringtopicString.format(agriculture/%s/command,deviceId);CommandMessagemsgnewCommandMessage();msg.setCommand(command);msg.setParams(params);msg.setTimestamp(System.currentTimeMillis());StringpayloadnewObjectMapper().writeValueAsString(msg);// 创建出站处理器MqttPahoMessageHandlerhandlernewMqttPahoMessageHandler(outboundClientId,mqttClientFactory);handler.setDefaultTopic(topic);handler.setDefaultQos(1);// 至少一次送达// 发送handler.handleMessage(MessageBuilder.withPayload(payload).build());}}注意出站适配器的clientId和入站的不能一样否则会冲突。这里加了-outbound后缀来区分。六、连接异常处理农业生产环境不如机房稳定MQTT 连接偶尔会断开。我们需要感知并处理这种状况BeanpublicMqttPahoClientFactorymqttClientFactory(){DefaultMqttPahoClientFactoryfactorynewDefaultMqttPahoClientFactory();MqttConnectOptionsoptionsmqttConnectOptions();// 方式一Paho 自带的自动重连推荐options.setAutomaticReconnect(true);// 重连间隔从 1 秒开始最大 30 秒options.setMaxReconnectDelay(30000);factory.setConnectionOptions(options);returnfactory;}/** * 方式二自定义回调监听连接状态更灵活 */ComponentpublicclassMqttConnectionListenerimplementsMqttCallbackExtended{OverridepublicvoidconnectComplete(booleanreconnect,StringserverURI){if(reconnect){// 重连成功处理离线期间积压的业务log.info(MQTT 重连成功: {},serverURI);sensorService.syncOfflineData();}else{log.info(MQTT 首次连接成功: {},serverURI);}}OverridepublicvoidconnectionLost(Throwablecause){log.error(MQTT 连接断开: {},cause.getMessage());// 可以在这里触发告警通知}OverridepublicvoidmessageArrived(Stringtopic,MqttMessagemessage){// 这个回调由 Paho 原生 API 触发// 用 Spring Integration 的话消息走 ServiceActivator这里不需要处理}}七、多 Topic 订阅与消息分发智慧农业场景下后端通常要同时订阅十几个 Topic。如果全堆在一个 Handler 里代码必然变成一锅粥。这时候策略模式就派上用场了// 定义处理器接口publicinterfaceTopicHandler{booleansupports(Stringtopic);voidhandle(Stringtopic,Stringpayload);}// 温度处理器ComponentpublicclassTemperatureHandlerimplementsTopicHandler{publicbooleansupports(Stringtopic){returntopic.contains(/temperature);}publicvoidhandle(Stringtopic,Stringpayload){// 温度相关处理}}// 湿度处理器ComponentpublicclassHumidityHandlerimplementsTopicHandler{publicbooleansupports(Stringtopic){returntopic.contains(/humidity);}publicvoidhandle(Stringtopic,Stringpayload){// 湿度相关处理}}// 统一分发器ComponentpublicclassMessageDispatcher{privatefinalListTopicHandlerhandlers;publicMessageDispatcher(ListTopicHandlerhandlers){this.handlershandlers;}publicvoiddispatch(Stringtopic,Stringpayload){for(TopicHandlerhandler:handlers){if(handler.supports(topic)){handler.handle(topic,payload);return;}}log.warn(未找到匹配的处理器: {},topic);}}Spring 会自动扫描所有实现了TopicHandler的 Bean注入到MessageDispatcher。后续新增传感器类型只需新增一个 Handler 类完全符合开闭原则——对扩展开放对修改关闭。总结SpringBoot 整合 MQTT 的关键步骤就三步配置连接参数 → 定义消息通道 → 绑定处理器。Spring Integration 帮你屏蔽了连接管理、线程调度、异常重试这些脏活累活你就可以专心写业务逻辑。记住几个容易踩的坑多实例部署时clientId必须唯一入站和出站不能用同一个clientIdsetCleanSession(false)配合setAutomaticReconnect(true)才是弱网环境的正确打开方式