去年我接手一个环境监测平台的后端改造,硬件端是几十台温湿度、空气质量采集网关,每台网关下面挂着若干传感器节点。旧方案是设备通过HTTP POST每30秒上报一次JSON,后端用SpringBoot接口接收再写入MySQL。设备规模到两百台左右时问题集中爆发:SIM卡网络抖动时HTTP重连风暴直接把Nginx连接数打满;要远程给设备下发校准指令,还得在网关写一套反向轮询逻辑,麻烦得要命。后来把设备接入全部切到MQTT,用SpringBoot做后端订阅端,从Topic订阅到传感器报文解析,这套架构稳定跑了大半年,今天把完整方案和踩过的坑一起整理出来。
这篇内容适合两类人:一类是刚接触物联网设备接入,想搞明白SpringBoot怎么对接MQTT的后端开发;另一类是已经在用MQTT,但被断线重连、报文解析、重复消费这些问题折磨过的同学。我不会从头讲SpringBoot基础,但MQTT协议里和工程落地直接相关的关键点会展开说清楚,保证看完能直接动手改造项目。
1. 设备接入选型:为什么是MQTT,而不是HTTP、TCP长连接或CoAP
1.1 传统HTTP上报的痛点
做设备接入,很多人第一个想到的就是HTTP接口。设备端定时POST数据,后端搞个Controller接收,看起来简单直接。但设备量大了之后,HTTP方案的短板非常明显。
首先是服务端无法主动下发指令。设备上报是单向的,平台要给设备下发校准指令、升级命令、远程控制,只能等设备下一次上报时在响应里捎带指令。如果设备上报频率是30秒一次,一条指令最坏要等30秒才能送达;如果设备上报频率是1小时一次,控制指令基本就废了。我之前那个项目需要远程校准设备时钟,HTTP方案下要么改设备端上报频率,要么只能等,体验很差。
其次是弱网环境下的连接风暴。工业网关大多用4G或Wi-Fi,网络质量不稳定。一旦信号抖动,几十台设备同时重连,Nginx的并发连接数瞬间飙升,后端服务也跟着抖。HTTP是短连接还好,但每次POST都要做TCP握手、TLS握手、HTTP头解析,开销全浪费在非业务部分。如果改成HTTP长连接轮询,服务端连接管理又是一堆事。
第三个问题是离线消息处理。设备断网期间上报的数据丢了就丢了,HTTP没有消息持久化的语义。业务方如果要求「设备离线期间数据不能丢」,HTTP方案你得自己搞队列补偿,而网络断的时长不可控,补偿逻辑很难写好。
1.2 MQTT、TCP长连接、CoAP三者的边界
既然HTTP不靠谱,那自己用Netty搞一个TCP长连接协议行不行?可以,但你要自己解决心跳机制、粘包拆包、断线重连、消息确认、会话恢复这些问题。每一样都是真实的工程成本。设备几百台时还能撑,到几千台上万台时,任何一个小坑都会被放大成大事故。
对比CoAP,它更适合极低功耗、受限网络环境下的设备,比如NB-IoT水表、路灯控制器。CoAP基于UDP,虽然轻量,但可靠传输要靠消息重发机制,实现复杂度其实不低,而且SpringBoot生态里对CoAP的成熟支持远不如MQTT丰富。
MQTT的定位正好卡在中间:比HTTP更适合长连接场景,比自研TCP协议省事得多,比CoAP在Java生态中更成熟。它专为低带宽、高延迟、弱网环境设计,发布订阅模型天然支持服务端下发指令,QoS机制给消息可靠性提供了明确分级。
做个表格直观对比一下:
| 维度 | HTTP短连接 | 自研TCP长连接 | CoAP | MQTT |
|---|---|---|---|---|
| 通信模型 | 请求/响应 | 自定义 | 请求/响应 + 观察 | 发布/订阅 |
| 服务端下发 | 不支持 | 自己设计 | 支持,机制简单 | 天然支持 |
| 弱网表现 | 差,重连风暴 | 靠代码质量 | 有重发机制 | QoS + 心跳 |
| 离线消息 | 需自己补偿 | 需自己实现 | 不支持 | Broker保留会话 |
| 开发成本 | 低 | 高 | 中 | 低 |
| Java生态 | 好 | 中 | 一般 | 好 |
从我实际体验来说,MQTT最大的优势不是某一个单点能力,而是把长连接场景下那些重复的脏活——心跳、重连、消息确认、会话恢复——全部标准化了。你只需要关注业务报文的格式和解析逻辑。
2. MQTT核心概念:从Broker、Topic到QoS,搞懂协议再动手
2.1 Broker、发布者、订阅者三个角色
MQTT协议里只有三个角色,比想象中简单。
Broker是消息服务器,负责接收所有消息并按主题转发。发布者是数据产生方,物联网场景里通常是传感器网关;订阅者是数据消费方,我们做的SpringBoot后端就扮演这个角色。一块数据从发布到被消费,中间不经过任何业务服务器,全部由Broker路由,这是MQTT和HTTP最大的思维差异——消息不是「发给谁」,而是「发到哪个主题」,谁关心这个主题谁就订阅它。
Broker的选型直接影响生产稳定性。我自己用过的方案有这几种:
- EMQX:开源版功能完整,支持共享订阅、插件扩展,控制台可观测性好,中小规模项目首选。
- Mosquitto:轻量级,适合内网测试或者嵌入式环境,单机性能一般,大规模生产不建议。
- HiveMQ:商业产品,高可用方案完善,适合对SLA要求极高的场景。
- 云厂商托管MQTT:像阿里云IoT、AWS IoT Core这类,免运维,但可能面临设备数据出网合规问题,需要评估。
如果只是开发联调,本地直接Docker起一个EMQX就够用。
docker run -d --name emqx -p 1883:1883 -p 18083:18083 emqx/emqx:5.8.01883是MQTT端口,18083是控制台端口,浏览器打开能实时看到连接数、消息量、订阅关系,排查问题非常方便。
2.2 Topic设计与通配符
Topic是MQTT的路由地址,用/分隔层级。设计Topic时要把设备信息、数据类型、业务维度放进去,类似设计URL路径。我常用的模式是:
factory/{工厂ID}/gateway/{网关ID}/telemetry也可以按设备类型分:
device/{deviceType}/{deviceId}/data订阅时可以用通配符一次订阅多个主题。+匹配单层,#匹配多层且只能在末尾。要接收所有网关的上报数据,写成factory/+/gateway/+/telemetry;要把某个网关下的所有数据都拿到,写成factory/001/gateway/+/#。
这里有两个容易被坑的点:第一,Topic不建议以/开头,某些Broker对以/开头的Topic处理不一致;第二,发布消息的Topic里不能带通配符,只有订阅才能用+和#,把+写进发布Topic里消息会直接发不出去。
2.3 QoS语义、cleanSession和遗嘱消息
QoS是MQTT协议里最核心的概念,它定义了消息投递的可靠性等级。
- QoS 0:最多一次。消息发出去了不管结果,可能丢失,适合低频状态上报、日志数据。
- QoS 1:至少一次。Broker收到消息后要回ACK,发送方没收到ACK就重发,消息可能重复,适合大多数传感器数据上报。
- QoS 2:恰好一次。通过四步握手保证不丢不重,性能开销最大,适合计费消息、控制指令这类不能错不能漏的场景。
设备上报传感器数据,默认建议用QoS 1。数据重复了可以靠业务层去重,但数据丢了不太好追,除非你允许漏采。
cleanSession这个参数也很关键。cleanSession=true时,设备断线后Broker清空它的会话状态,离线消息不保留;cleanSession=false时,Broker会保存设备的订阅关系和离线消息,设备重新上线后能收到断线期间的消息。代价是Broker要持续维护会话状态,对服务端内存有消耗。
还有两个容易忽略但很实用的机制:遗嘱消息和保留消息。设备连接时可以在Broker上登记一条遗嘱Topic,当设备异常掉线时,Broker立刻往这个Topic发一条消息,方便后端感知设备离线。保留消息则是让Broker为某个Topic保留最后一条消息,新订阅者订阅后立刻收到,适合下发设备配置这种场景。
3. SpringBoot整合MQTT:客户端配置与连接管理实操
3.1 依赖选型:Paho还是Spring Integration MQTT
SpringBoot整合MQTT,有大方向上的两种选择:直接用Eclipse Paho客户端,或者用Spring Integration MQTT模块。
Spring Integration MQTT封装程度高,把MQTT抽象成消息通道,和Spring的@MessageEndpoint、@ServiceActivator等注解配合得很顺。如果你只做简单的收发,用它非常快。但问题在于,它把底层连接生命周期管理封装得比较死,遇到需要精细控制重连、自定义回调、订阅恢复的场景,反而束手束脚。
Eclipse Paho是MQTT协议的官方Java客户端实现,API直白,MqttClient、MqttConnectOptions、MqttCallback这些类一看就懂,回调、重连、会话恢复都是原生支持的。生产环境里我倾向直接上Paho,自己控制连接管理,心里有底。
引入依赖:
<dependency> <groupId>org.eclipse.paho</groupId> <artifactId>org.eclipse.paho.client.mqttv3</artifactId> <version>1.2.5</version> </dependency>Paho同时提供了同步API的MqttClient和异步API的MqttAsyncClient。实际用下来,MqttClient的subscribe()是阻塞等待SUBACK的,业务上更容易确认订阅结果,我就以它为例。
3.2 连接配置:从application.yml到MqttConnectOptions
配置文件里把Broker地址、账号密码、clientId、心跳间隔放好:
mqtt: broker-url: tcp://localhost:1883 client-id: backend-service-01 username: backend password: secret connection-timeout: 10 keep-alive-interval: 30 clean-session: false topic: gateway/+/telemetryclientId在同一个Broker上必须全局唯一,客户端和Broker都靠它标识会话。如果两个连接用了同一个clientId,后连上来的会把前面那个踢下线,这是生产环境最经典的故障之一。
连接配置类:
@Configuration public class MqttConfig { @Value("${mqtt.broker-url}") private String brokerUrl; @Value("${mqtt.client-id}") private String clientId; @Value("${mqtt.username}") private String username; @Value("${mqtt.password}") private String password; @Value("${mqtt.connection-timeout}") private int connectionTimeout; @Value("${mqtt.keep-alive-interval}") private int keepAliveInterval; @Bean(destroyMethod = "close") public MqttClient mqttClient() throws MqttException { MqttConnectOptions options = new MqttConnectOptions(); options.setAutomaticReconnect(true); options.setCleanSession(false); options.setConnectionTimeout(connectionTimeout); options.setKeepAliveInterval(keepAliveInterval); options.setUserName(username); options.setPassword(password.toCharArray()); options.setMaxInflight(200); MqttClient client = new MqttClient(brokerUrl, clientId); client.connect(options); client.setCallback(new MqttCallbackHandler(client)); return client; } }setAutomaticReconnect(true)是Paho提供的自动重连机制,Broker暂时不可用时会按策略自动重连。setMaxInflight(200)设置了未确认消息的最大数量,如果设备量大、消息密集,这个值太小会导致发送端消息积压。
这里有个启动时序的坑:应用启动时Broker可能还没就绪,如果@Bean里直接connect()会抛异常导致Spring容器启动失败。生产环境我的做法是把连接动作放到ApplicationReadyEvent之后触发,或者用@Retryable做几次重试,避免硬件环境没准备好直接把后端服务搞挂。
3.3 订阅与回调:把报文送进业务线程池
订阅的核心是消息回调。Paho的MqttCallback有四个方法:connectionLost连接断开时触发、messageArrived收到消息时触发、deliveryComplete消息发送完成后触发,以及扩展回调MqttCallbackExtended里的connectComplete在连接成功时触发。
我建议用MqttCallbackExtended,因为自动重连成功后需要重新订阅,只有它能在连接建立时给你一个钩子:
@Component public class MqttCallbackHandler implements MqttCallbackExtended { private static final Logger log = LoggerFactory.getLogger(MqttCallbackHandler.class); private final MqttClient mqttClient; private final SensorDataParser parser; private final DeviceDataService deviceDataService; private final ExecutorService executor; public MqttCallbackHandler(MqttClient mqttClient, SensorDataParser parser, DeviceDataService deviceDataService) { this.mqttClient = mqttClient; this.parser = parser; this.deviceDataService = deviceDataService; this.executor = new ThreadPoolExecutor( 4, 8, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(2000), new ThreadPoolExecutor.CallerRunsPolicy()); } @Override public void connectComplete(boolean reconnect, String serverURI) { log.info("MQTT连接建立完成,reconnect={}, serverURI={}", reconnect, serverURI); subscribeTopics(); } @Override public void connectionLost(Throwable cause) { log.warn("MQTT连接断开,等待自动重连", cause); } @Override public void messageArrived(String topic, MqttMessage message) { // 回调线程是Paho内部线程,不能在这里做耗时操作 executor.submit(() -> processMessage(topic, message.getPayload())); } private void processMessage(String topic, byte[] payload) { try { SensorData sensorData = parser.parse(payload); deviceDataService.process(topic, sensorData); } catch (Exception e) { log.error("报文处理失败,topic={}", topic, e); } } private void subscribeTopics() { try { String[] topics = {"gateway/+/telemetry", "gateway/+/event"}; int[] qos = {1, 1}; mqttClient.subscribe(topics, qos); log.info("MQTT订阅成功"); } catch (MqttException e) { log.error("MQTT订阅失败", e); } } }这里关键点是messageArrived里绝对不要直接做数据库写入或者报文解析这种耗时操作。Paho的回调线程是单线程的,如果你在回调里做慢SQL,后续所有消息都会被阻塞,消息积压越来越严重。正确做法是丢给业务线程池。
线程池的拒绝策略我用了CallerRunsPolicy,意思是队列满了之后,回调用当前线程处理。这在高峰期可以避免消息丢失,只是回调线程会被业务拖慢,后续消息进来会排队,但这比丢数据要好。
4. 传感器报文解析:从字节到业务数据的完整链路
4.1 一张典型的传感器数据帧长什么样
MQTT消息的Payload本身只是二进制字节流,设备厂商一般会定义自己的报文协议。我接触过的传感器报文,绝大多数是自定义的二进制帧格式,JSON传输反而少见,因为二进制体积小、解析快,适合MCU级别的嵌入式设备。
拿一个典型的环境监测传感器举例,它的数据帧格式这样定义:
| 字段 | 长度(字节) | 说明 |
|---|---|---|
| 帧头 | 2 | 固定 0xAA 0x55,用于同步 |
| 设备类型 | 1 | 0x01 环境监测,0x02 电力监测 |
| 设备ID | 4 | 小端字节序,例如 0x01000000 表示设备ID=1 |
| 数据类型 | 1 | 0x01 周期上报,0x02 告警上报 |
| 数据区长度 | 1 | 数据区的字节数N |
| 数据区 | N | 传感器数据体 |
| CRC16 | 2 | 从「设备类型」到「数据区末尾」计算得到的CRC16 |
| 帧尾 | 2 | 固定 0x0D 0x0A |
数据区里面,环境监测类型的结构是定长的:
| 字段 | 长度(字节) | 类型/单位 |
|---|---|---|
| 温度 | 2 | 有符号short,除以10,单位℃ |
| 湿度 | 2 | 无符号short,除以10,单位%RH |
| PM2.5 | 2 | 无符号short,单位μg/m³ |
| 噪声 | 2 | 无符号short,除以10,单位dB |
| 电池电压 | 2 | 无符号short,单位mV |
定长结构的报文最好解析,因为字段偏移量固定。实际设备里还会遇到数据区里有多个记录块的情况,比如设备一次性上报10分钟内每秒一条的采样数据,本质就是从数据区里循环读取定长结构,原理完全相同。
4.2 用Java解析二进制报文:ByteBuffer、字节序与位运算
解析这类报文,Java的ByteBuffer是最顺手的工具。它有getShort()、getInt()、get()等方法,配合ByteOrder.LITTLE_ENDIAN处理小端字节序,省去手动移位拼字节的麻烦。
解析的核心步骤是:校验帧头帧尾、校验CRC、按偏移量读取各字段。
public SensorData parse(byte[] payload) { if (payload.length < 13) { throw new IllegalStateException("报文长度不足,实际长度=" + payload.length); } // 1. 校验帧头帧尾 if (payload[0] != (byte) 0xAA || payload[1] != (byte) 0x55) { throw new IllegalStateException("帧头校验失败"); } if (payload[payload.length - 2] != (byte) 0x0D || payload[payload.length - 1] != (byte) 0x0A) { throw new IllegalStateException("帧尾校验失败"); } // 2. CRC校验 int crcOffset = payload.length - 4; int crcExpected = Crc16Util.calc(payload, 2, crcOffset); int crcActual = (payload[crcOffset] & 0xFF) | ((payload[crcOffset + 1] & 0xFF) << 8); if (crcExpected != crcActual) { throw new IllegalStateException("CRC校验失败,期望=" + crcExpected + ",实际=" + crcActual); } // 3. 按字段读取 ByteBuffer buf = ByteBuffer.wrap(payload).order(ByteOrder.LITTLE_ENDIAN); buf.position(2); // 跳过帧头 byte deviceType = buf.get(); int deviceId = buf.getInt(); byte dataType = buf.get(); int dataLen = buf.get() & 0xFF; if (buf.remaining() < dataLen + 4) { throw new IllegalStateException("数据区长度与报文实际长度不符"); } byte[] dataArea = new byte[dataLen]; buf.get(dataArea); ByteBuffer dataBuf = ByteBuffer.wrap(dataArea).order(ByteOrder.LITTLE_ENDIAN); float temperature = dataBuf.getShort() / 10.0f; float humidity = Short.toUnsignedInt(dataBuf.getShort()) / 10.0f; int pm25 = Short.toUnsignedInt(dataBuf.getShort()); float noise = Short.toUnsignedInt(dataBuf.getShort()) / 10.0f; int voltageMv = Short.toUnsignedInt(dataBuf.getShort()); SensorData data = new SensorData(); data.setDeviceId(deviceId); data.setDeviceType(deviceType); data.setDataType(dataType); data.setTemperature(temperature); data.setHumidity(humidity); data.setPm25(pm25); data.setNoise(noise); data.setVoltageMv(voltageMv); data.setTimestamp(System.currentTimeMillis()); return data; }注意几个容易写错的地方:
buf.getShort()返回的是有符号short,如果温度是负数,直接用没问题;但湿度、PM2.5这种只有非负值的字段,高位是1时直接强转int会变成负数,必须用Short.toUnsignedInt()转成无符号整数。同理,CRC计算时每个字节都要& 0xFF,否则会被符号位干扰。
字节序务必和设备端确认清楚。大多数MCU是小端存储,但少数设备用大端。一旦端序搞反,解析出的数据全是乱的。拿到报文样例后,先在纸上按字节手算一遍,确认端序再写代码,不要凭感觉。
4.3 数据不完整、脏数据、多记录交错怎么办
MQTT协议本身在传输层保证了每条消息的字节流是完整的,不会出现TCP那类半包粘包问题。协议头里的剩余长度字段已经给出了消息边界,Paho会帮你把一条消息的所有字节组装好再交给回调。但应用层仍然要面对三类问题:
第一类是数据区长度字段和设备真实数据不一致。有些硬件工程师写固件不严谨,数据区长度算错,多一个字节少一个字节的情况都有。解析代码里要做长度校验,数据区长度超过Payload剩余长度的直接丢弃并告警,不要硬解析,否则读出来的字段全是错位的。
第二类是CRC校验失败。设备在通信干扰下可能发出错帧。CRC校验是最后一道防线,它对数据区做完整性验证,不通过的消息必须丢掉,不能入库,不然会出现一个温度-45℃、湿度140%的脏数据,你的统计报表直接废掉。
第三类是设备上报频率高,回调线程处理不过来。这在前面线程池那里已经说了,队列配合合理配置能顶住一阵子,但设备量持续涨的话,还是要靠水平扩展——多部署几个后端实例,用MQTT的共享订阅把消息分散到不同实例处理,而不是靠单机硬扛。
5. 生产环境避坑:连接不稳定、重复消息、报文丢失怎么办
5.1 自动重连后订阅丢失:你掉进了cleanSession的陷阱
Paho的setAutomaticReconnect(true)让客户端在断线后能自动重连,但很多人不知道,自动重连成功后,订阅关系不一定会恢复。
原因在于cleanSession的设置。cleanSession=false时,Broker会保存设备的订阅关系,重连后自动恢复;cleanSession=true时,Broker不保存任何会话状态,重连后你的订阅关系是空的,如果不重新订阅,消息永远不会到达。
我的配置里写的是cleanSession=false,但这不代表万事大吉。Broker重启、会话过期、clientId变动,都可能导致会话丢失。所以最保险的做法仍然是在connectComplete回调里无条件重新订阅,幂等操作,多调一次没损失,少调一次就丢数据。
5.2 QoS1的消息重复消费:要有幂等设计
QoS1的语义是「至少一次」,Broker在收不到ACK时会重发消息,这意味着你的消息处理逻辑必须容忍重复。
重复的来源不止QoS重发,还有设备端自身的行为——很多设备固件在重启后会重新上报上一次的数据,或者网络恢复后把积压的消息重新发一遍。这些都是正常现象,不能靠设备端解决,只能靠业务端做幂等。
我常用的方案是在报文协议里加入消息序号,数据区前面带一个4字节的自增序号,后端用「设备ID + 序号」作为唯一键。写入时先查RedisSETNX:
String dedupKey = "dedup:msg:" + deviceId + ":" + seq; Boolean firstSeen = redisTemplate.opsForValue().setIfAbsent(dedupKey, "1", Duration.ofMinutes(30)); if (Boolean.FALSE.equals(firstSeen)) { // 重复消息,丢弃 log.info("重复报文已忽略,deviceId={}, seq={}", deviceId, seq); return; }如果不用Redis,用数据库唯一索引也行,原理一样。但千万不要以为设备量小就不会重复,我遇到过一台网关一晚上重复上报了6000条数据,没有幂等机制的话,MySQL直接爆掉。
5.3 clientId冲突:多实例部署时的定时炸弹
前面提过clientId在Broker上必须唯一。如果后端服务部署了多个实例,每个实例用相同的clientId连接同一个Broker,结果就是后面连接的实例把前面的踢下线,前面的重连再踢后面的,两个实例反复互踢,日志里全是连接断开的告警,消息消费基本瘫痪。
这问题看起来蠢,但非常容易踩。解决办法是让clientId带上实例标识,比如:
String clientId = "backend-" + UUID.randomUUID().toString().substring(0, 8);或者用微服务里已有的实例ID、Pod名称来生成,确保每次启动都不同。注意如果用cleanSession=false,clientId变动会导致之前的会话丢失,所以多实例部署时cleanSession=false基本没意义,共享会话只有单实例场景能用。
多实例消费的场景,正确做法是改回cleanSession=true,配合MQTT的共享订阅,让Broker把消息负载均衡到多个实例上。
5.4 回调线程被阻塞:慢SQL拖垮整个消息链路
Paho的messageArrived回调运行在客户端内部线程上,这个线程不仅处理消息回调,还负责发送心跳PINGREQ。如果你在回调里直接执行慢SQL或者调用第三方接口,一个慢请求卡住几秒,会导致客户端没空闲处理心跳,Broker判定设备超时,把连接断开,然后触发重连,重连后又要重新订阅,整个系统进入不稳定状态。
我的处理思路一直是把消息处理和连接管理彻底隔离。messageArrived只做一件事——把消息丢进业务线程池。线程池参数根据设备量和单条消息处理的耗时来算。假设一台网关每秒上报10条消息,100台网关就是每秒1000条,单条处理耗时20ms,那需要的线程数是1000 * 0.02 = 20个线程。我一开始配的4个线程明显不够,后来改成核心线程数16、最大线程数32才算稳。
线程池队列长度也要给出上限。我见过队列配成无界队列,高峰期积压几十万条消息,内存直接飙到OOM的案例。队列满的时候用CallerRunsPolicy或者记录丢弃日志,都比内存爆炸要好。
5.5 消费成功但入库失败:丢失的不是消息,是业务数据
MQTT的ACK机制发生在协议层:客户端收到消息后,Paho会自动向Broker返回PUBACK,Broker就认为消息投递成功。如果你的业务代码在入库时失败了,此时ACK已经发出去了,Broker不会重发这条消息。
这是MQTT使用中最容易产生误解的地方。协议层的「送达」只是说Broker把消息交到了客户端手里,不代表业务处理成功。
我的做法是解析入库不在消息抵达时同步做,而是先把报文存到一个异常重试表,状态标记为「已接收待处理」。入库成功就更新状态为「已完成」,入库失败就定时任务重试。如果确认是脏数据不可修复,人工标记丢弃。
这套机制的好处是,无论代码出bug还是数据库抖动,消息都不会丢,最多延迟处理。对传感器数据来说,延迟几秒钟根本不影响业务,但丢了就是不可逆的损失。
6. 消息确认与可观测性:让问题在爆发前暴露
6.1 从连接和消费两个维度做健康度量
MQTT接入跑起来之后,最难的不是功能开发,而是出了问题怎么快速定位。我建议至少埋四个指标:
- 连接状态:当前是否连接、断线次数、自动重连次数。
- 消息量:每秒接收的消息数,按Topic维度统计。
- 处理耗时:从收到消息到处理完成的耗时分布,P99尤其重要。
- 失败率:CRC校验失败数、解析失败数、入库失败数。
连接状态通过MqttCallbackExtended里的connectComplete和connectionLost很好统计,加一个计数器就行:
private final AtomicInteger connectCount = new AtomicInteger(); private final AtomicInteger disconnectCount = new AtomicInteger(); @Override public void connectComplete(boolean reconnect, String serverURI) { if (reconnect) { reconnectCount.incrementAndGet(); } else { connectCount.incrementAndGet(); } } @Override public void connectionLost(Throwable cause) { disconnectCount.incrementAndGet(); }消息量、耗时、失败率建议用Micrometer直接接Prometheus,一套标准方案,图省事也可以先打印成日志再用日志采集工具处理。我发现很多故障其实早有预兆,只是没人在意这些指标。
设备侧的掉线检测,靠的是Broker的心跳超时和遗嘱消息。设备连接时设置遗嘱Topic为gateway/{id}/online,遗嘱内容为offline,正常上线时发一条online。后端订阅这个Topic,就能实时感知设备上下线。心跳间隔设备端设多少,后端要做兼容,有的设备设30秒,有的设60秒,Broker超时时间要留足余量,一般是心跳间隔的1.5倍。
6.2 处理失败的补偿链路:别让一条脏数据阻塞后续处理
解析失败的报文不要放在内存里就完了,把原始字节持久化下来。我通常建一张mqtt_failed_message表,字段包括:Topic、原始报文hex、失败原因、失败时间、重试次数。每天定时扫一遍,统计失败率趋势。如果某类设备持续CRC失败,基本可以判定是硬件固件问题,把hex拿出来找设备厂商对线,一张表搞定。
入库失败的消息单独一张message_retry表,后台任务每30秒重试一次,连续失败超过5次标记为阻塞,人工介入。这里有个经验:重试的任务要控制并发,别几百条失败消息一次性全查出来重试,把数据库打挂。
这里还要注意一条:如果你的消息里带了时间戳,尽量用设备上报的时间而不是服务端接收时间。设备延迟上报的数据,用服务端时间会导致统计准点率出错。传感器报文里一般都有采样时间字段,没有的话宁可让设备端加一个,也不要自己脑补。
6.3 关于共享订阅和水平扩展的补充
单实例处理能力总有上限。当设备规模到K级别,每秒消息量上万条时,单后端实例哪怕线程池配置再合理,也会出现瓶颈。
MQTT的解法是共享订阅。订阅时把Topic写成$share/{group}/{topic},比如$share/backend/gateway/+/telemetry。同一个group里的多个订阅者,Broker会把消息轮询分发,每个消息只发给其中一个。这样后端服务只需要水平扩展实例数,加一台实例就多一份消费能力,完全不用改业务代码。
这是我后来在设备量翻倍时用的方案。切共享订阅要注意:如果当前用的是cleanSession=false加共享订阅,某些Broker实现下行为会有差异。EMQX 4.x/5.x对共享订阅支持都好,但老版本Mosquitto对共享订阅的支持有坑。切换前先在测试环境验证一下。
多Topic想要不同分组做不同处理,比如上报数据和告警事件分开处理,就订阅两个组:
$share/data-group/gateway/+/telemetry $share/event-group/gateway/+/event业务上把两类消息消费完全隔离,告警处理线程池和数据处理线程池分开,即使数据量大导致数据处理堵了,告警通道还是畅通的——设备出问题的时候,告警反而是你最需要的信号。
关于这套架构的一点后续想法
整个SpringBoot整合MQTT的链路,从选型、连接管理到报文解析、生产避坑,到这里已经完整了。能跑通Demo不难,难的是在设备量增长、网络抖动、固件不靠谱的情况下还能稳定运行。
我另外两个补充建议:一是设备端的MQTT SDK和协议文档,一定要在项目启动前对齐,字节序、字段类型、CRC算法这些细节确认得越早,后续返工越少;二是测试环境尽量模拟弱网,用tc netem加一些丢包和延迟再跑联调,这比生产环境半夜被叫醒处理故障划算太多了。
如果你正在改造类似的设备接入系统,可以按这篇文章的思路先搭骨架,再根据你们设备的具体协议补充细节。有拿不准的报文格式或者连接问题,欢迎留言交流。