
简介这份资源面向在SpringBoot环境下开发物联网通信的Java工程师聚焦eclipse.paho.client.mqttv3客户端的集成与生产级改造。内容覆盖MQTT连接建立、消息订阅发布、断线重连与心跳检测、线程池高并发处理以及消息落库MySQL与缓存Redis的完整业务流程适合需要将MQTT通信从Demo推进到稳定可用的中高级开发者。资源包共128个文件以103个xml配置与16个java源码为主另含yml、sql、md说明及Windows端mosquitto安装程序整体约25.61MB目录结构清晰便于按模块查阅与二次改造。目前已有197人学习下载。读者可直接获得可运行的客户端示例、重连与线程池改造思路、数据库与缓存写入代码以及配套配置与建表脚本快速在本地搭建测试环境并迁移到实际项目。1. 从一次 MQTT 消息积压说起SpringBoot 集成 paho.mqttv3 到底要解决什么设备侧每秒推 300 条状态上报SpringBoot 服务端用默认的 MQTT 回调直接写库跑了不到两小时MySQL 连接池被打满Redis 写入延迟飙到 800ms客户端还被 broker 踢下线——这是我在一个物联网项目里真实遇到的场景。问题不在 MQTT 协议本身而在于eclipse.paho.client.mqttv3的MqttCallback回调线程模型所有消息都在 Paho 内部的单线程里串行回调你在messageArrived里做任何耗时操作写库、调接口、发 Redis都会阻塞后续消息的接收最终触发 broker 的 keepAlive 超时断连。这篇要讲清楚的就是在 SpringBoot 里用eclipse.paho.client.mqttv3搭一个能扛住高并发、断线能自己爬回来、消息可靠落到 MySQL 和 Redis 的 MQTT 客户端。核心要解决四件事——连接生命周期管理、断线重连策略、回调线程池改造、消息存储入库。适合正在做设备接入、消息中间件对接、或者被 MQTT 消息积压折磨过的后端同学。下面按「先跑通最小连接 → 再改线程池 → 再做存储 → 最后排坑」的顺序推。2. 最小可运行paho.mqttv3 客户端接入 SpringBoot 的完整配置2.1 依赖引入与 MqttClient 还是 MqttAsyncClient 的选型Paho 的 Java 客户端有两个核心类MqttClient是同步阻塞的MqttAsyncClient是全异步的。很多人一上来就用MqttClient然后在connect()那里卡住主线程。我的建议是服务端接入一律用MqttAsyncClient它的connect()返回IMqttToken可以注册IMqttActionListener做回调不会阻塞 SpringBoot 的启动流程。Maven 依赖只需要一个dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version /dependency版本选 1.2.5 是因为它在 1.2.x 系列里对重连的automaticReconnect支持最稳定1.2.0 之前有几个重连后 session 丢失的已知问题。如果你用的是 SpringBoot 2.7注意它默认带的spring-integration-mqtt会传递依赖一个旧版本需要显式排除掉否则会出现NoSuchMethodError。2.2 连接参数配置cleanSession、keepAlive、automaticReconnect 三个必调项连接选项是踩坑最密集的地方直接上配置类Configuration public class MqttConfig { Value(${mqtt.broker}) private String broker; Value(${mqtt.clientId}) private String clientId; Value(${mqtt.username}) private String username; Value(${mqtt.password}) private String password; Bean public MqttAsyncClient mqttAsyncClient() throws MqttException { // 内存持久化避免重启后未确认消息丢失生产可换 MqttDefaultFilePersistence MemoryPersistence persistence new MemoryPersistence(); MqttAsyncClient client new MqttAsyncClient(broker, clientId, persistence); MqttConnectOptions options new MqttConnectOptions(); options.setUserName(username); options.setPassword(password.toCharArray()); options.setCleanSession(false); // 关键false 才能保留会话和离线消息 options.setKeepAliveInterval(30); // 30s小于 broker 的 1.5 倍超时阈值 options.setConnectionTimeout(10); // 连接超时 10s options.setAutomaticReconnect(true); // 开启自动重连 options.setMaxInflight(100); // 未确认消息上限默认 10 太小 client.connect(options).waitForCompletion(10000); return client; } }cleanSessionfalse是断线重连能收到离线消息的前提broker 会为这个 clientId 保留 session 和 QoS0 的未确认消息。但要注意clientId 必须全局唯一且固定如果你用 UUID 每次生成session 永远对不上。keepAliveInterval设 30 秒broker 端一般按 1.5 倍即 45 秒判定超时留出网络抖动余量。maxInflight默认只有 10高并发下消息会堵在客户端队列里发不出去调到 100 是常见做法。2.3 订阅与回调为什么不能在 messageArrived 里直接写库订阅本身很简单client.subscribe(device//status, 1); // QoS 1至少一次QoS 选择上设备状态上报用 1 就够QoS 2 的四次握手在高并发下开销翻倍除非是计费类消息否则没必要。真正的问题在回调client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { // 重连成功后必须重新订阅cleanSessionfalse 时 broker 会恢复但显式订阅更保险 try { client.subscribe(device//status, 1); } catch (MqttException e) { log.error(重新订阅失败, e); } } Override public void messageArrived(String topic, MqttMessage message) { // 这里如果直接 jdbcTemplate.update(...)就是灾难的开始 String payload new String(message.getPayload(), StandardCharsets.UTF_8); // 丢给线程池立刻返回 messageExecutor.submit(() - processMessage(topic, payload)); } Override public void connectionLost(Throwable cause) { log.warn(MQTT 连接断开: {}, cause.getMessage()); } Override public void deliveryComplete(IMqttDeliveryToken token) { } });messageArrived运行在 Paho 的CommsCallback单线程里你在这里每多花 1ms后面排队的消息就多等 1ms。正确做法是只做「解析 topic 投递线程池」两件事业务逻辑全部异步化。connectComplete回调是MqttCallbackExtended才有的普通MqttCallback没有重连后重新订阅必须靠它这是很多人重连后收不到消息的根因。3. 线程池高并发改造从单线程回调到 ThreadPoolExecutor 的落地细节3.1 回调线程池的参数怎么定核心数、队列、拒绝策略Paho 回调是单线程我们要在回调里把消息转交给自己的线程池。这个池子的参数不能拍脑袋得按消息吞吐量算。假设峰值 5000 条/秒单条处理耗时 20ms含写库那需要的并发线程数约等于5000 × 0.02 100。但线程不是越多越好MySQL 连接池通常也就 50-100线程数超过连接池只会让线程在getConnection()上排队。我一般这样配Bean(messageExecutor) public ThreadPoolExecutor messageExecutor() { int coreSize Runtime.getRuntime().availableProcessors() * 2; return new ThreadPoolExecutor( coreSize, // 核心线程数 coreSize * 2, // 最大线程数 60L, TimeUnit.SECONDS, // 空闲回收时间 new LinkedBlockingQueue(2000), // 有界队列防止 OOM new ThreadFactoryBuilder().setNameFormat(mqtt-msg-%d).build(), new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略 ); }队列选LinkedBlockingQueue而不是SynchronousQueue是因为 MQTT 消息允许短暂缓冲SynchronousQueue在突发流量下会直接触发拒绝。队列容量 2000 是个经验值太小容易丢消息太大则内存压力和延迟都上去了。拒绝策略用CallerRunsPolicy而不是AbortPolicy是因为 MQTT 消息丢了就真丢了让调用线程也就是 Paho 回调线程自己跑相当于给 broker 一个背压信号——回调变慢Paho 收消息变慢broker 那边 inflight 满了自然降速。这是用阻塞换可靠性的经典取舍。3.2 消息处理链路解析、幂等、批量入库的拆分线程池只是把并发打开了真正决定吞吐的是处理链路。我的做法是把一条消息拆成三段解析、幂等判断、入库。private void processMessage(String topic, String payload) { // 1. 解析topic 形如 device/{deviceId}/status String[] parts topic.split(/); String deviceId parts[1]; DeviceStatus status JSON.parseObject(payload, DeviceStatus.class); // 2. 幂等用 msgId deviceId 做 Redis SETNX防止 QoS1 重复投递 String idempotentKey mqtt:idem: status.getMsgId(); Boolean first redisTemplate.opsForValue() .setIfAbsent(idempotentKey, 1, 5, TimeUnit.MINUTES); if (Boolean.FALSE.equals(first)) { return; // 重复消息直接丢弃 } // 3. 入库先写 Redis 缓存最新状态再异步落 MySQL redisTemplate.opsForHash().putAll(device:status: deviceId, status.toMap()); mysqlQueue.add(new DeviceStatusRecord(deviceId, status)); }QoS 1 是「至少一次」意味着同一条消息可能到达多次幂等必须做。用 Redis 的SETNX加过期时间是最轻量的方案key 里带 msgId5 分钟过期足够覆盖重投窗口。Redis 存最新状态用 Hash 结构一个设备一个 key方便前端直接查。MySQL 落库不要在这里同步做再丢一个队列做批量插入每 500 条或每 1 秒 flush 一次能把 MySQL 的写入压力降一个数量级。3.3 背压与限流当线程池队列满了怎么办队列满了触发CallerRunsPolicyPaho 回调线程被占用这其实是好事——它形成了天然的背压。但你要监控这个信号否则等到 broker 把客户端踢了才发现。加一个定时任务打印线程池状态Scheduled(fixedRate 10000) public void monitorPool() { ThreadPoolExecutor pool (ThreadPoolExecutor) messageExecutor; log.info(MQTT线程池 active{}, queue{}, completed{}, rejected{}, pool.getActiveCount(), pool.getQueue().size(), pool.getCompletedTaskCount(), rejectedCount.get()); if (pool.getQueue().size() 1500) { log.warn(MQTT消息队列积压当前 {}考虑扩容消费者, pool.getQueue().size()); } }队列持续超过容量的 75% 就是扩容信号。扩容方向有两个加线程受限于 MySQL 连接池或者加消费者实例水平扩展。如果单机线程已经到连接池上限就该考虑把消息先写 Kafka 再消费而不是硬扛。4. 存储入库MySQL 与 Redis 的分工和写入策略4.1 Redis 该存什么最新状态、去重标记、设备在线心跳Redis 在这个链路里承担三个角色别混用用途数据结构Key 示例过期策略设备最新状态Hashdevice:status:{id}不设过期靠业务清理消息幂等去重Stringmqtt:idem:{msgId}5 分钟设备在线心跳Stringdevice:online:{id}90 秒靠续期维持最新状态用 Hash 是因为设备字段多Hash 可以单字段更新不用整体覆盖。心跳用 String 加过期时间设备每次上报就SET一次刷新 TTL90 秒没上报就自动消失前端查EXISTS就知道设备是否在线。这里有个坑不要用device:status:*这种通配 key 做扫描KEYS命令在生产环境会阻塞 Redis要用SCAN或者维护一个设备 ID 的 Set。4.2 MySQL 批量入库攒批、事务、失败重试MySQL 这边用批量插入攒批逻辑Component public class MysqlBatchWriter { private final ListDeviceStatusRecord buffer new ArrayList(500); private final Object lock new Object(); Scheduled(fixedDelay 1000) public void flush() { ListDeviceStatusRecord batch; synchronized (lock) { if (buffer.isEmpty()) return; batch new ArrayList(buffer); buffer.clear(); } try { jdbcTemplate.batchUpdate( INSERT INTO device_status(device_id, msg_id, status, ts) VALUES(?,?,?,?), batch, 500, (ps, record) - { ps.setString(1, record.getDeviceId()); ps.setString(2, record.getMsgId()); ps.setInt(3, record.getStatus()); ps.setLong(4, record.getTs()); }); } catch (Exception e) { log.error(批量入库失败{} 条消息丢失, batch.size(), e); // 失败重试写本地文件或重新入队别直接吞掉 } } public void add(DeviceStatusRecord record) { synchronized (lock) { buffer.add(record); if (buffer.size() 500) { flush(); } } } }batchUpdate的 batchSize 设 500配合rewriteBatchedStatementstrue的 JDBC 参数能把 500 条单插变成一条多值 INSERT性能差 10 倍以上。失败重试不要简单重试因为可能是某条数据格式问题导致整批失败我的做法是失败后拆成单条逐条插把坏数据挑出来单独记录。4.3 事务边界MQTT 消息的「至少一次」和数据库「恰好一次」怎么对齐MQTT 保证的是至少一次数据库想要的是恰好一次中间靠幂等对齐。顺序很重要先写 Redis 幂等标记再写 MySQL。如果反过来MySQL 写成功但 Redis 标记失败重投时就会重复入库。Redis 标记先写即使后续 MySQL 失败重投时会被幂等拦住——代价是这条消息丢了但至少不会脏数据。要更严格的话用 MySQL 的唯一索引兜底ALTER TABLE device_status ADD UNIQUE KEY uk_msg (msg_id);插入时用INSERT IGNORE或ON DUPLICATE KEY UPDATE让数据库层做最终去重。这样即使 Redis 挂了也不会重复。5. 避坑与排查断线重连、线程池、存储链路的 5 个血泪教训5.1 重连后收不到消息cleanSession 和 clientId 的坑现象网络恢复后日志显示connectComplete回调触发但订阅的 topic 一条消息都收不到。原因两个可能。一是cleanSessiontruebroker 每次连接都当新会话之前的订阅全丢二是 clientId 用了随机值重连时 broker 认为是新客户端旧 session 被顶掉。解决cleanSession必须设falseclientId 用「服务名 固定后缀」的格式比如iot-server-node1保证每次重连都是同一个身份。另外在connectComplete里显式重新subscribe一次双保险。5.2 线程池把内存打爆无界队列的隐形杀手现象服务运行几小时后 OOM堆 dump 里全是LinkedBlockingQueue$Node。原因用了new LinkedBlockingQueue()无参构造队列容量是Integer.MAX_VALUE消息积压时无限堆积。解决队列必须设容量上限2000 是常用值。同时加监控队列超过 75% 就告警。如果业务确实需要大缓冲用ArrayBlockingQueue预分配内存比LinkedBlockingQueue的节点对象开销小。5.3 Redis 连接超时导致消息处理线程全挂现象Redis 网络抖动几秒之后 MQTT 消息处理全部卡死线程池 active 数打满。原因processMessage里 Redis 操作没有超时设置默认阻塞直到 TCP 超时可能几十秒线程池线程全被占住。解决Redis 客户端配置连接超时和读写超时Lettuce 设spring.redis.timeout2000msJedis 设connectionTimeout和soTimeout。同时给 Redis 操作加熔断连续失败就降级为直接写 MySQL别让一个依赖拖垮整条链路。5.4 QoS 1 重复消息导致数据翻倍现象MySQL 里同一设备的同一 msgId 出现多条记录。原因QoS 1 的「至少一次」语义broker 没收到 PUBACK 就会重投客户端重启或网络抖动都会触发。解决三层防护——Redis SETNX 幂等标记、MySQL 唯一索引、插入用INSERT IGNORE。三层里任何一层生效都能挡住重复别只靠一层。5.5 keepAlive 设置过大导致断线检测迟钝现象设备实际已经离线但服务端 5 分钟后才感知到。原因keepAliveInterval设了 300 秒broker 要等 1.5 倍即 450 秒才判定超时。解决keepAlive 设 30-60 秒配合客户端的心跳包。代价是心跳流量增加但断线检测从分钟级降到秒级。如果设备侧省电优先可以适当放宽但服务端要有独立的在线状态超时逻辑别完全依赖 MQTT 的 keepAlive。6. 进阶技巧用 MqttCallbackExtended 做重连后的状态补偿前面讲的都是「消息来了怎么处理」但有个场景容易被忽略断线期间设备上报的消息重连后 broker 会补发QoS0 且 cleanSessionfalse可如果断线时间很长补发的消息可能已经过期直接入库会污染最新状态。我的做法是在connectComplete里做一次状态补偿。Override public void connectComplete(boolean reconnect, String serverURI) { if (reconnect) { log.info(MQTT 重连成功开始补偿设备状态); // 1. 重新订阅 try { client.subscribe(device//status, 1); } catch (MqttException e) { log.error(重订阅失败, e); } // 2. 拉取断线期间的最新状态覆盖可能过期的补发消息 ListDevice devices deviceMapper.selectAll(); for (Device device : devices) { String cached redisTemplate.opsForValue().get(device:online: device.getId()); if (cached null) { // 缓存已过期说明设备断线期间没上报标记为离线 redisTemplate.opsForHash().put(device:status: device.getId(), online, 0); } } } }这段逻辑的核心是重连后不要盲目相信补发的消息先用 Redis 里的在线状态做一次对账。如果device:online:{id}已经过期说明设备在断线期间没有心跳那补发的历史消息就不应该覆盖当前状态直接标记离线更准确。另一个技巧是给消息加时间戳入库时判断msgTs是否小于 Redis 里记录的最后更新时间小于就丢弃。这样即使补发消息乱序到达也不会把新状态覆盖成旧状态。Long lastTs (Long) redisTemplate.opsForHash().get(device:status: deviceId, ts); if (lastTs ! null status.getTs() lastTs) { return; // 过期消息丢弃 }这个「时间戳水位线」的做法本质上是用 Redis 维护一个每设备的最新时间戳所有到达的消息都要过这道闸。代价是每次处理多一次 Redis 读但相比数据错乱带来的排查成本这点开销完全值得。我自己在这个项目上最大的教训是MQTT 客户端的可靠性不取决于你用了多高级的 API而取决于你对「至少一次」语义的敬畏。一开始我觉得 QoS 1 已经够可靠了没做幂等结果上线第一周就出现数据翻倍排查了两天才定位到是 broker 重投。后来把幂等、唯一索引、时间戳水位线三层都加上才真正睡了个安稳觉。线程池参数也别一次调到位先按公式算个初值上线后看监控慢慢调队列积压和拒绝计数是最诚实的指标。希望帮到你。本文还有配套的精品资源点击获取