发布时间:2026/9/1 0:50:46
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)才是弱网环境的正确打开方式

相关新闻

2026/9/1 0:50:46

前端工程原型如何补齐稳定性边界

前端工程原型如何补齐稳定性边界将 Vue3 响应式系统接入大模型流式 SSE(Server-Sent Events)推送时,若更新频率过高,页面可能出现卡顿。 即使只是持续追加文本,也可能让输入和滚动变慢。是否由响应式更新引起&#xff…

2026/9/1 0:45:45

研发流程自动化

研发流程自动化先确认范围 流程自动化先从构建、测试和发布记录的衔接处开始,不急着把所有决策交给模型。 记录可复查的依据 每个步骤都应声明输入、产物、失败处理和负责人;版本变化时,旧记录只能参考,不能直接复用结论。 用小范…

2026/9/1 0:45:45

Django实战:开发停车场预约计费系统的完整指南

简介:本资源是一套基于Python Django框架开发的停车场预约与计费系统完整源码案例,面向Web开发初学者及Django进阶学习者,聚焦真实业务场景中的用户管理、车位调度、动态计费与在线支付等核心功能实现。压缩包共2000个文件,含1632…

2026/9/1 1:05:46

Stable Diffusion漫画助手工作流:从安装到批量出图的完整指南

简介:本资源是面向数字绘画创作者与AI绘画初学者的Stable Diffusion漫画辅助工具安装与使用指南,聚焦解决漫画角色设计、构图优化与色彩辅助等高频创作痛点,依托人工智能技术提升出图效率与艺术表现力。压缩包共5个文件(929KB&…

2026/9/1 1:05:46

Python+YOLOv8打造智能驾驶员状态监测系统:从训练到UI实现

简介:本资源是一套基于Python与YOLOv8实现的智能驾驶员状态监测系统完整项目,面向高校毕业设计、课程设计及AI视觉开发初学者,聚焦疲劳驾驶行为检测这一典型工业落地场景。项目支持闭眼、张嘴、睁眼、闭嘴四类关键状态识别,含约30…

2026/9/1 1:05:46

AI风险工程化实践:大模型应用防护层与治理体系搭建

AI技术正在以极快的速度进入生产系统,但近两年科技界关于 AI 风险的讨论也变得越来越频繁。不管是科技高管在公开场合表达担忧,还是企业内部对模型失控、信息污染、隐私泄露的讨论,本质上都指向同一个问题:大模型能做什么只是能力…

2026/9/1 1:05:46

从线稿协作到AI辅助上色:指绘接力工作流完整复盘

指绘“大妖精接力”工作流复盘:从线稿协作到AI辅助上色的完整技术方案这一篇不聊空泛的绘画理念,直接拆解一个具体项目:8月8日进行的“大妖精接力雾湖上妖精”。这个项目名字看起来是粉丝创作活动,但背后藏的是一套非常典型的多人…

2026/9/1 1:05:46

AI绘画角色一致性出图:ComfyUI+LoRA+ControlNet批量生成同人图工作流

这是很多二次元同人圈里常说的一句话。围绕“卡莫大人”这个虚拟角色,粉丝、画师、创作者会集中产出同人图、表情包和贺图。如果把这句话放到技术语境里,其实刚好对应一个很实际的需求:如何用 AI 绘画给同一个虚拟角色批量生成风格统一、角度…

2026/9/1 1:00:46

网络服务的安全核算

网络服务的安全核算先确定问题 网络服务的安全核算的讨论先落在授权范围、样本来源和保存方式。不要用一段笼统的经验替代前提:输入从哪里来、谁负责确认、失败后怎样停止,都应在开始前写清。 沿着一条路径检查 围绕网络服务的安全核算做安全研究实践时&…

2026/8/31 1:05:20

vSound小提琴数字处理器实操指南:从接线到演出的完整配置

电小提琴或者原声小提琴插电演出,第一个绕不开的坎就是声音难听。原声琴的共鸣和空气感一旦进了拾音器,出来的往往是一坨干瘪、发尖、带着奇怪塑料味的信号。我当初第一次把琴接上乐队调音台,直接被主唱吐槽"你这声音像在锯钢丝"。…

2026/8/31 2:14:20

传感器接口IC如何攻克生物化学传感的微弱信号难题?

1. 从电极到比特流:为什么生物化学传感必须依赖专用接口IC 做生物化学传感的人都有过类似的经历:明明传感器本身性能很好,信号输出却一塌糊涂——噪声大、漂移明显、重复性差,怎么调都达不到预期。很多时候问题并不在传感器&#…

2026/8/31 1:41:28

STM32F411CEU6多通道ADC采集:扫描模式+DMA实现详解

1. 多通道 ADC 的用武之地把“Multichannel ADC”和“STM32F411CEU6”这两个关键字放在一起,其实就是嵌入式开发里最常遇到的一类需求:用一块不算贵的 MCU,同时采集多路模拟信号。STM32F411CEU6 是 48 引脚的 Cortex-M4F 主控,主频…

2026/9/1 0:00:42

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/1 0:00:42

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/1 0:00:42

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

2026/9/1 0:00:42

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行&#xff0c;Type-C接口算是典型的“看着简单&#xff0c;做起来全坑”的东西。光引脚就24个&#xff0c;高低速信号、电源、控制线全部塞在一个小小的连接器里&#xff0c;如果PCB布局不做规划&#xff0c;打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/1 0:00:42

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/1 0:00:42

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…