拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SpringBoot整合MQTT实现软硬件通信实战

SpringBoot整合MQTT实现软硬件通信实战 1. 项目概述为什么软硬件通信必须跨过“协议鸿沟”在工业现场、智能楼宇、农业物联网这些真实场景里我见过太多团队卡在同一个地方后端服务写得再漂亮前端页面再炫酷一到要跟温湿度传感器、PLC控制器、电表采集器、LoRa网关这些硬件设备“说上话”整个系统就哑火了。不是数据收不到就是指令发不出或者偶尔通一下第二天又断连。问题往往不在于SpringBoot代码写错了而在于——你根本没意识到软件世界和硬件世界之间横着一道“协议鸿沟”。这道鸿沟的本质是通信范式的彻底错位。我们写Java Web应用默认用HTTP靠RESTful接口、JSON格式、无状态请求响应模型而绝大多数嵌入式设备、工业控制器、低功耗终端它们的通信逻辑是轻量、异步、事件驱动的讲究的是“有消息就推过来”而不是“你来问我有没有新数据”。MQTT协议就是专门为填平这道鸿沟而生的。它不是HTTP的替代品而是为资源受限设备和不可靠网络量身定制的“消息邮局”发布者把信投进指定邮箱Topic订阅者只要提前登记好邮箱地址信到了就立刻收到中间不依赖双方同时在线也不需要复杂的握手流程。所以“SpringBoot整合MQTT实现软硬件通信”这个标题拆开来看核心不是“怎么配个依赖”而是解决三个层次的问题第一层让SpringBoot这个企业级Java框架能像呼吸一样自然地收发MQTT消息而不是硬塞一个阻塞式客户端进去第二层把硬件端发来的原始字节流比如Modbus RTU解析后的16进制报文映射成业务可理解的Java对象比如TemperatureReading第三层设计出健壮的消息路由与错误兜底机制——当网关掉线5分钟传感器还在疯狂上报消息不能丢也不能压垮内存。我去年帮一家水表厂商做远程抄表系统就因为没处理好第三层凌晨三点被运维电话叫醒发现MQTT Broker内存爆满所有未确认消息堆积如山。后来我们加了QoS 1本地磁盘缓存死信队列重试才真正稳住。这不是配置问题是架构意识问题。这个项目适合三类人一是正在做IoT平台、能源监控、设备管理系统的Java后端开发者你需要把设备数据真正接入业务流二是刚从传统Web开发转向物联网方向的工程师需要补上“设备侧思维”这一课三是硬件工程师或系统集成商想快速验证自己设备的MQTT通信是否符合标准。它不教你怎么写单片机代码但会告诉你当你把0x01 0x03 0x00 0x01 0x00 0x02 0xC4 0x0B这样的Modbus帧发到device/001A2B/reading这个Topic时SpringBoot后端该怎么接、怎么转、怎么存、怎么告警。接下来的内容全部围绕这四个动作展开——接、转、存、告。2. 整体架构设计与技术选型逻辑2.1 为什么不是HTTP、WebSocket或自定义TCP直面协议选型的硬伤很多团队一开始会本能地想“既然SpringBoot天生支持HTTP那让硬件设备也走HTTP POST不就行了”我试过也踩过坑。去年给一个冷链运输车队做温控系统第一批方案就是让车载终端每30秒POST一次JSON到/api/v1/temperature。上线两周后运维报警API网关CPU持续95%日志里全是Connection reset by peer。查下来不是代码问题是协议本身不匹配。HTTP是“请求-响应”模型每次通信都要建立TCP连接、TLS握手、发送Header、等待响应对一个电池供电、GPRS带宽只有20Kbps的车载终端来说光握手就占了80%的流量和电量。更致命的是HTTP没有“推送”能力——你想让服务器主动下发一条“立即上传全量数据”的指令只能靠终端轮询这又进一步加剧了网络负担。WebSocket看起来是个折中方案它支持双向通信连接复用。但问题在于它的“长连接”特性在物联网场景下反而成了负担。WebSocket连接需要服务器维持大量空闲连接状态每个连接至少占用几KB内存。当你的设备规模从100台涨到10万台服务器内存直接翻百倍而其中99%的连接99%的时间都在“挂机”。MQTT的“发布/订阅”模型则完全不同Broker只负责消息路由不维护客户端状态。一个MQTT连接可以承载成千上万个Topic的订阅连接本身极轻量心跳包只有2字节且天然支持QoS服务质量分级——你可以对关键告警消息设QoS 2确保送达对普通心跳设QoS 0尽力而为这种灵活性是HTTP和WebSocket无法提供的。至于自定义TCP协议听起来很“硬核”但代价极高。你需要自己实现连接保活、消息分包粘包、重传机制、加密协商、心跳超时检测……这些轮子早被MQTT协议族打磨了二十年。Eclipse Paho、EMQX、Mosquitto这些成熟组件背后是全球数千名工程师对各种网络异常弱网、断电、DNS漂移的反复锤炼。我们做的是业务系统不是协议栈研发。选择MQTT本质是选择站在巨人的肩膀上把精力聚焦在“温度数据来了之后要不要触发告警工单”这种业务逻辑上而不是“怎么让TCP包不丢”。2.2 SpringBoot整合MQTT的三种路径为什么最终锁定spring-integration-mqtt在Spring生态里整合MQTT有三条主流路径我逐一实测对比过路径一原生Paho Client Scheduled轮询这是最“原始”的方式。引入org.eclipse.paho:org.eclipse.paho.client.mqttv3手写MqttClient连接、订阅、消息回调。优点是完全可控缺点是“反Spring”——你需要手动管理连接生命周期、线程池、异常重连所有Bean都得自己new。更麻烦的是消息到达后如何把它变成Spring容器里的事件你得自己发ApplicationEvent再写EventListener监听。这套流程写下来50行代码里有30行是胶水代码而且一旦连接断开重连逻辑极易出错。我见过有团队在重连时没清空旧的MqttCallback导致同一条消息被重复消费两次。路径二spring-boot-starter-integrationspring-integration-mqtt这是Spring官方推荐的集成方式。它把MQTT客户端封装成Spring Integration的MessageChannel和MessageHandler消息进来自动转成Message?对象你可以用ServiceActivator注解轻松绑定处理方法。最关键的是它深度融入Spring生命周期连接由MqttPahoClientFactory管理自动重连、自动恢复订阅、自动处理QoS。你甚至可以用InboundChannelAdapter声明一个“消息源”就像注入一个DataSource一样自然。配置项也极其清晰mqtt.url、mqtt.username、mqtt.password、mqtt.default.qos全部通过application.yml控制。我们线上系统用的就是这条路稳定运行18个月零MQTT连接相关故障。路径三Spring Cloud Stream Binder for MQTT这是面向云原生的方案把MQTT当作一种“消息中间件抽象”通过spring-cloud-stream-binder-mqtt让你的代码完全不感知MQTT只写StreamListener。理论上很美但实际落地时发现两个硬伤一是Binder社区维护较弱最新版只支持到Spring Boot 2.7对3.x支持滞后二是它过度抽象当你要调试某个Topic的QoS级别或自定义Will Message遗嘱消息时得层层穿透Binder源码远不如路径二直接。对于大多数企业级IoT项目路径二的平衡性最好足够简单足够强大足够稳定。所以最终技术栈锁定为Spring Boot 3.2.x spring-integration-mqtt 6.2.x EMQX 5.7作为MQTT Broker。EMQX选5.7是因为它对QoS 2的支持最完善且内置规则引擎后续可以无缝对接数据库存储避免我们在SpringBoot里写一堆JDBC代码。2.3 架构分层图从物理设备到业务告警的完整链路整个通信链路不是简单的“设备→MQTT→SpringBoot”而是一个多层过滤与转换的流水线。我画了一个简化的分层图文字描述它决定了后续所有代码的设计逻辑[物理层] 温湿度传感器 → Modbus RTU → 485转WiFi网关 → MQTT Publish (Topic: device/{mac}/raw) [协议层] 网关固件将Modbus响应帧如01 03 02 14 05 B9 2F解析为JSON再发布到MQTT {type:reading,voltage:5.2,temp:23.5,humi:45.1,ts:1715823456} [传输层] EMQX Broker接收消息根据Topic路由执行预设规则如若temp50则转发到告警Topic [应用层 - SpringBoot] 1. Inbound Channel Adapter监听 device//raw 2. 消息经JsonToObjectTransformer转为DeviceRawData对象 3. Service Activator调用DeviceDataProcessor.process() → 校验数据有效性如温度范围-40~85℃ → 转换为业务实体TemperatureReading → 写入TimescaleDB时序数据库 → 若触发阈值发布告警消息到 alarm/device/{mac} 4. Outbound Channel Adapter监听 alarm/#推送微信/短信通知这个分层的核心思想是每一层只做一件事且这件事必须做到极致。物理层专注信号稳定协议层专注格式统一传输层专注路由可靠应用层专注业务逻辑。很多项目失败就是因为试图在应用层同时干四件事——既要解析十六进制又要校验CRC又要存库又要发邮件结果一处出错全链路雪崩。而按这个分层如果某天要换掉EMQX换成HiveMQ只需改Broker配置如果网关固件升级JSON格式变了只需调整JsonToObjectTransformer的映射规则业务处理器DeviceDataProcessor一行代码都不用动。3. 核心细节解析与实操要点3.1 MQTT Broker选型与本地化部署为什么EMQX比Mosquitto更适合生产环境Broker是整个MQTT通信的“心脏”选错它后面所有优化都是空中楼阁。市面上最常被提及的是Mosquitto和EMQX很多人觉得Mosquitto轻量、开源、文档全就直接上生产。我必须坦白在中小规模1万设备的验证环境Mosquitto确实够用但一旦进入真实产线它的短板就会暴露无遗。Mosquitto最大的问题是缺乏企业级运维能力。它没有内置的Dashboard所有监控都得靠mosquitto_sub命令行工具抓取统计Topic如$SYS/broker/clients/connected然后自己写脚本解析。当凌晨设备批量掉线你得在终端里敲十几条命令才能定位是网络问题还是认证失败。更麻烦的是集群——Mosquitto官方集群方案基于mosquitto.conf的cluster配置但实际部署时节点间同步延迟高经常出现“消息在一个节点收到另一个节点收不到”的情况这对需要强一致性的告警系统是灾难。EMQX则完全不同。它从设计之初就面向云原生5.x版本采用eMQL语言规则引擎把消息路由、数据清洗、协议转换全部可视化配置。举个真实例子我们水表系统要求所有上报的meter/001A2B/data消息必须先检查battery_level字段是否低于20%如果是就自动转发到alarm/battery_lowTopic并附带设备位置信息。在Mosquitto里这得写一个Python脚本监听Topic解析JSON判断条件再调用Paho Client发布新消息——四步操作三个单点故障。在EMQX里一条SQL就能搞定SELECT clientid as device_id, payload.battery_level as battery, payload.location as position FROM meter//data WHERE payload.battery_level 20然后把这条SQL的输出目标设为alarm/battery_low。规则引擎会自动编译成Erlang字节码在Broker内核里执行毫秒级延迟零额外进程。本地部署EMQX我推荐Docker方式因为它能完美规避Windows/Linux/macOS的环境差异。以下是经过生产验证的docker-compose.yml精简版version: 3.8 services: emqx: image: emqx/emqx:5.7.0 container_name: emqx ports: - 1883:1883 # MQTT TCP - 8081:8081 # Dashboard - 8883:8883 # MQTT SSL environment: - EMQX_NAMEemqx - EMQX_HOSTnode1.emqx.io - EMQX_CLUSTER__DISCOVERYstatic - EMQX_CLUSTER__STATIC__SEEDSemqxnode1.emqx.io - EMQX_LOADED_PLUGINSemqx_management,emqx_recon,emqx_retainer,emqx_rule_engine volumes: - ./emqx_data:/opt/emqx/data - ./emqx_log:/opt/emqx/log - ./emqx_conf/emqx.conf:/opt/emqx/etc/emqx.conf restart: unless-stopped关键点在于volumes映射emqx.conf必须自定义禁用默认的匿名登录allow_anonymous false并配置JWT鉴权authentication [ { mechanism jwt, from password, key your-super-secret-jwt-key-2024, algorithm HS256 } ]这样设备连接时用户名填设备MAC地址密码填用该密钥签发的JWT Token既安全又免去了维护用户表的麻烦。我测试过单节点EMQX 5.7在4核8G服务器上能稳定支撑3万并发连接消息吞吐达12000 QPS完全满足绝大多数工业场景。3.2 SpringBoot端MQTT配置的魔鬼细节QoS、Clean Session与Will MessageSpringBoot整合MQTT的配置看似简单但几个关键参数的取值直接决定系统是“稳如老狗”还是“三天两头掉线”。我在application.yml里写了整整一页注释这里提炼出最易被忽视的三点第一QoS服务质量不是越高越好QoS有0、1、2三级。QoS 0是“最多一次”消息发出去就不管了QoS 1是“至少一次”Broker会存一份副本等客户端ACK才删除QoS 2是“恰好一次”通过四次握手确保不重不漏。很多教程一上来就教设QoS 2这是大忌。QoS 2的握手开销是QoS 0的4倍在GPRS或NB-IoT网络下一次QoS 2的发布可能耗时3秒以上严重拖慢设备响应。我们的实践是设备上报数据用QoS 1允许少量重复但绝不能丢失服务器下发指令用QoS 2必须确保设备收到。配置如下spring: integration: mqtt: inbound: default-qos: 1 # 所有订阅默认QoS 1 outbound: default-qos: 2 # 所有发布默认QoS 2第二Clean Session必须设为false这是导致“设备上线收不到历史消息”的元凶。clean-sessiontrue默认值意味着每次连接Broker都会丢弃该客户端之前的所有会话状态包括未ACK的消息、订阅关系。设备重启后它相当于一个全新客户端自然收不到断连期间别人发给它的消息。正确做法是设为false并配合client-id使用唯一标识如设备MACBean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory new DefaultMqttPahoClientFactory(); factory.setServerURIs(new String[]{tcp://localhost:1883}); factory.setUserName(admin); factory.setPassword(public.getBytes()); // 关键禁用clean session启用持久会话 factory.setConnectionOptionsCustomizer(options - { options.setCleanSession(false); options.setAutomaticReconnect(true); options.setKeepAliveInterval(60); // 心跳间隔60秒 }); return factory; }这样即使设备断网8小时重新连上Broker所有QoS0的未确认消息都会被重新推送。第三Will Message遗嘱消息是最后的安全网当设备异常断开如断电、看门狗复位Broker会自动发布一条预设的“遗嘱消息”到指定Topic通知系统“设备离线了”。这比心跳超时检测快得多。配置方法是在连接选项里设置options.setWill(device/status, offline.getBytes(), 1, true);意思是如果连接意外中断Broker自动向device/statusTopic发布offline消息QoS为1且消息保留Retained。这样任何新订阅device/status的客户端一上来就能看到最新的设备状态无需再查数据库。3.3 消息编解码从原始字节流到业务对象的精准映射硬件设备发来的消息99%不是标准JSON而是二进制帧或自定义文本格式。比如我们合作的某款电表上报数据是这样的十六进制字符串01 03 06 00 00 00 01 00 00 7D 0A。如果直接当字符串处理SpringBoot收到的就是乱码。必须在消息进入业务逻辑前完成“协议解析”这一步。我们采用Spring Integration的Transformer链式处理分三步走第一步HexStringToByteArrayTransformer先写一个自定义Transformer把十六进制字符串转成byte数组Component public class HexStringToByteArrayTransformer implements GenericTransformerString, byte[] { Override public byte[] transform(String source) { if (source null || source.trim().isEmpty()) { return new byte[0]; } // 移除空格按两位分割 String hex source.replaceAll(\\s, ); byte[] bytes new byte[hex.length() / 2]; for (int i 0; i bytes.length; i) { bytes[i] (byte) Integer.parseInt(hex.substring(i * 2, i * 2 2), 16); } return bytes; } }然后在Integration Flow里引用Bean public IntegrationFlow mqttInboundFlow(MqttPahoClientFactory factory) { return IntegrationFlow.from( Mqtt.messageDrivenChannelAdapter(factory) .uri(tcp://localhost:1883) .clientid(springboot-server) .topic(device//raw) .qos(1) ) .transform(Transformers.fromJson(DeviceRawData.class)) // 如果是JSON直接转 .transform(hexStringToByteArrayTransformer) // 如果是Hex先转byte[] .transform(new ModbusFrameParser()) // 自定义解析器提取寄存器值 .channel(c - c.executor(Executors.newFixedThreadPool(10))) .get(); }第二步ModbusFrameParser这是一个专门解析Modbus RTU帧的类它知道帧结构[Address][Function][Data][CRC]。我们只关心Data段比如上面的例子00 00 00 01 00 00代表4个寄存器值分别是0, 1, 0, 0。解析后封装成ModbusResponse对象public class ModbusResponse { private int address; private int function; private ListInteger registers; // 解析出的寄存器值列表 // getter/setter... }第三步业务对象组装最后把ModbusResponse转成业务实体ElectricMeterDataServiceActivator(inputChannel mqttInboundChannel) public ElectricMeterData processModbusResponse(ModbusResponse response) { ElectricMeterData data new ElectricMeterData(); data.setDeviceId(001A2B); // 从Topic里提取 data.setVoltage(response.getRegisters().get(0) * 0.1); // 寄存器0是电压单位0.1V data.setCurrent(response.getRegisters().get(1)); // 寄存器1是电流单位mA data.setEnergy(response.getRegisters().get(2) response.getRegisters().get(3) * 65536L); // 32位电能高低位组合 data.setTimestamp(System.currentTimeMillis()); return data; }这个三层解析链的好处是职责单一可测试性强。你可以单独给ModbusFrameParser喂一个byte数组断言它返回的registers列表是否正确也可以模拟一个ModbusResponse测试processModbusResponse的业务逻辑。比起把所有解析逻辑堆在ServiceActivator方法里这种写法在设备协议变更时修改成本低一个数量级。4. 实操过程与核心环节实现4.1 从零搭建SpringBoot MQTT服务5分钟可运行的最小可行代码很多教程一上来就贴几百行配置新手看得头皮发麻。我给你一个绝对能跑通的“最小可行代码”MVP只有4个文件5分钟内就能在本地验证通路。记住目标不是写完美系统而是先让“设备发消息SpringBoot打印出来”这件事发生。第一步创建SpringBoot项目用Spring Initializrhttps://start.spring.io/选以下依赖Spring WebSpring IntegrationSpring Integration MQTTLombok简化getter/setter第二步添加MQTT依赖在pom.xml里确认有dependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-integration/artifactId /dependency第三步配置application.ymlspring: integration: mqtt: inbound: default-qos: 1 outbound: default-qos: 1 mqtt: url: tcp://localhost:1883 username: admin password: public client-id: springboot-server # 日志调成DEBUG方便看MQTT连接过程 logging: level: org.springframework.integration.mqtt: DEBUG org.eclipse.paho: DEBUG第四步编写核心配置类创建MqttConfig.javaConfiguration EnableIntegration public class MqttConfig { Value(${spring.mqtt.url}) private String mqttUrl; Value(${spring.mqtt.username}) private String username; Value(${spring.mqtt.password}) private String password; Value(${spring.mqtt.client-id}) private String clientId; Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory new DefaultMqttPahoClientFactory(); factory.setServerURIs(new String[]{mqttUrl}); factory.setUserName(username); factory.setPassword(password.getBytes()); factory.setConnectionOptionsCustomizer(options - { options.setCleanSession(false); options.setAutomaticReconnect(true); options.setKeepAliveInterval(60); }); return factory; } Bean public MessageChannel mqttInputChannel() { return MessageChannels.publishSubscribe().get(); } Bean public IntegrationFlow mqttInboundFlow(MqttPahoClientFactory factory) { return IntegrationFlow.from( Mqtt.messageDrivenChannelAdapter(factory) .uri(mqttUrl) .clientid(clientId) .topic(test/#) // 订阅test开头的所有Topic .qos(1) ) .channel(c - c.messageChannel(mqttInputChannel())) .get(); } ServiceActivator(inputChannel mqttInputChannel) public void handleMessage(Message? message) { String topic message.getHeaders().get(mqtt_receivedTopic, String.class); Object payload message.getPayload(); System.out.println(【收到消息】Topic: topic , Payload: payload); // 这里就是你的业务处理入口 } }第五步启动并测试先启动EMQXDocker命令docker-compose up -d启动SpringBoot应用观察控制台你会看到类似Connected to tcp://localhost:1883的日志用MQTT.fx工具免费桌面客户端连接localhost:1883用户名admin密码public发布一条消息Topic填test/helloPayload填{msg:Hello from MQTT.fx}QoS选1切回SpringBoot控制台你应该立刻看到打印【收到消息】Topic: test/hello, Payload: {msg:Hello from MQTT.fx}恭喜你已经打通了MQTT通信的第一公里。后续所有复杂功能——数据解析、存库、告警——都是在这个MVP基础上叠加的。记住这个原则永远先让最简路径跑通再逐步增强。我见过太多团队一上来就想做JWT鉴权、QoS 2、集群部署结果连本地连接都连不上白白浪费三天。4.2 设备端模拟用Python脚本代替真实硬件进行全流程联调没有真实硬件怎么验证整个链路别急用Python写一个超轻量的MQTT设备模拟器50行代码搞定。它能模拟设备上线、定时上报、异常断线、指令响应等所有关键行为比买硬件便宜比等硬件到货快。安装依赖pip install paho-mqtt创建device_simulator.pyimport json import time import random import paho.mqtt.client as mqtt from datetime import datetime # 设备配置 DEVICE_ID simulator-001 MQTT_BROKER localhost MQTT_PORT 1883 MQTT_USER admin MQTT_PASS public def on_connect(client, userdata, flags, rc): if rc 0: print(f[{DEVICE_ID}] 连接成功开始上报数据...) # 连接成功后发布上线消息 client.publish(fdevice/{DEVICE_ID}/status, online, qos1, retainTrue) else: print(f[{DEVICE_ID}] 连接失败返回码: {rc}) def on_disconnect(client, userdata, rc): print(f[{DEVICE_ID}] 意外断开连接返回码: {rc}) def simulate_sensor_data(): 模拟传感器数据生成 return { device_id: DEVICE_ID, temperature: round(random.uniform(20, 30), 1), humidity: random.randint(40, 70), battery: round(random.uniform(3.0, 4.2), 2), timestamp: int(datetime.now().timestamp() * 1000) } if __name__ __main__: client mqtt.Client(client_idfdevice-{DEVICE_ID}) client.username_pw_set(MQTT_USER, MQTT_PASS) client.on_connect on_connect client.on_disconnect on_disconnect # 连接Broker client.connect(MQTT_BROKER, MQTT_PORT, 60) client.loop_start() # 启动后台循环 try: while True: # 模拟每5秒上报一次 data simulate_sensor_data() topic fdevice/{DEVICE_ID}/reading client.publish(topic, json.dumps(data), qos1) print(f[{DEVICE_ID}] 已发布: {topic} - {data}) # 模拟1%概率异常断线测试Will Message if random.random() 0.01: print(f[{DEVICE_ID}] 模拟异常断线...) client.disconnect() time.sleep(3) client.reconnect() print(f[{DEVICE_ID}] 重新连接...) time.sleep(5) except KeyboardInterrupt: print(f\n[{DEVICE_ID}] 模拟器已停止) client.disconnect() client.loop_stop()运行这个脚本它会自动连接EMQX发布device/simulator-001/status为onlineretain消息每5秒生成一条模拟温湿度数据发布到device/simulator-001/reading有1%概率模拟设备异常断电触发Broker的Will Message你会在SpringBoot控制台看到offline消息支持CtrlC优雅退出现在把前面的SpringBootMqttConfig里的订阅Topic改成device//reading重启应用。你就能在控制台实时看到模拟器发来的每一条JSON数据。这个脚本的价值在于它把“硬件不可控”的问题转化成了“代码可控”的问题。你可以随时修改simulate_sensor_data()函数生成极端数据如温度-50℃、电池0.1V测试你的业务逻辑是否健壮也可以注释掉client.publish()只留on_disconnect专门测试断线重连逻辑。这才是高效联调的正确姿势。4.3 业务逻辑落地从消息到数据库再到告警的完整闭环现在消息能收到了模拟器也跑起来了下一步就是把“收到消息”变成“产生业务价值”。我们以一个真实的水表远程抄表场景为例实现从原始数据到微信告警的完整闭环。需求梳理水表每小时上报一次读数total_water单位m³和电池电量battery当battery 2.5V时触发低电量告警当total_water相比上次上报增长超过100m³疑似漏水触发异常用水告警所有数据存入PostgreSQL并建立时序索引第一步定义数据模型Data Builder NoArgsConstructor AllArgsConstructor public class WaterMeterReading { private String deviceId; private Long timestamp; // 毫秒时间戳 private Double totalWater; // 累计用水量 private Double battery; // 电池电压 private Double lastTotalWater; // 上次累计用水量用于计算增量 } // 告警实体 Data Builder public class AlarmEvent { private String deviceId; private String type; // battery_low, water_leak private String message; private Long triggerTime; }第二步编写业务处理器创建WaterMeterProcessor.javaService Slf4j public class WaterMeterProcessor { // 用ConcurrentHashMap缓存设备上次读数避免频繁查库 private final MapString, WaterMeterReading lastReadings new ConcurrentHashMap(); ServiceActivator(inputChannel mqttInputChannel) public void processWaterReading(Message? message) { try { // 1. 解析JSON为对象 String payload (String) message.getPayload(); JSONObject json new JSONObject(payload); String deviceId extractDeviceIdFromTopic(message); WaterMeterReading current WaterMeterReading.builder() .deviceId(deviceId) .timestamp(json.optLong(timestamp, System.currentTimeMillis())) .totalWater(json.optDouble(total_water, 0.0)) .battery(json.optDouble(battery, 0.0)) .build(); // 2. 检查低电量 if (current.getBattery() 2.5) { sendAlarm(AlarmEvent.builder() .deviceId(deviceId) .type(battery_low) .message(电池电压过低 current.getBattery() V) .triggerTime(current.getTimestamp()) .build()); } // 3. 检查异常用水需获取上次读数 WaterMeterReading last lastReadings.get(deviceId); if (last ! null current.getTotalWater() 0) { double increment current.getTotalWater() - last.getTotalWater(); if (increment 100.0) { sendAlarm(AlarmEvent.builder() .deviceId(deviceId) .type(water_leak) .message(疑似漏水1小时内用水量 increment m³) .triggerTime(current.getTimestamp()) .build()); } } // 4. 更新缓存并存库 lastReadings.put(deviceId, current); saveToDatabase(current); } catch (Exception e) { log.error(处理水表数据失败, e); } } private String extractDeviceIdFromTopic(Message? message) { String topic message.getHeaders().get(mqtt_receivedTopic, String.class); // Topic格式device/{deviceId}/reading提取deviceId return topic.split(/)[1]; } private void sendAlarm(AlarmEvent alarm) { // 发布告警消息到alarm/{deviceId} Topic由另一个Flow监听并推送微信 messagingTemplate.convertAndSend(alarm/ alarm.getDeviceId(), alarm); log.info(已触发告警{}, alarm); } private void saveToDatabase(WaterMeterReading reading) { // 使用JdbcTemplate插入此处省略具体SQL // INSERT INTO water_meter_readings (device_id, timestamp, total_water, battery) VALUES (?, ?, ?, ?) } }第三步配置告警推送Flow在MqttConfig.java里追加Bean public IntegrationFlow alarmOutboundFlow(MqttPaho
返回列表