ARTICLE DETAIL

资讯详情

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

SpringBoot集成MQTT高并发稳定接入实战

SpringBoot集成MQTT高并发稳定接入实战 简介MQTT是一种轻量级物联网消息协议广泛应用于智能硬件、工业IoT和车联网等设备直连场景。其核心原理基于发布/订阅模型与QoS分级保障机制技术价值在于低带宽占用、弱网适应性强及支持海量终端接入。典型应用场景包括设备状态上报、远程指令下发与实时告警推送。然而原生Paho客户端在SpringBoot中直接使用时极易遭遇断线重连失效、线程阻塞、MySQL写入瓶颈及Redis缓存不一致等生产级问题。本文聚焦MQTT客户端高并发改造深入解析断线重连状态机设计与线程池分级调度策略结合MySQL强一致写入和Redis最终一致性缓存双写方案提供可落地的稳定性工程实践。1. 项目概述为什么一个MQTT客户端要折腾这么多事SpringBoot里集成MQTT表面看就是加个依赖、写个连接、收发几条消息——但真放到生产环境跑一周你就会发现断网重连失败、消息堆积卡死、数据库写满磁盘、Redis连接池耗尽、线程数飙到200还压不住并发……这些不是“可能遇到的问题”而是所有没做过高并发MQTT接入的团队踩坑的必经之路。我去年在做工业IoT平台时就用eclipse.paho.client.mqttv3搭了第一版客户端。当时只想着“能连上broker就行”结果上线第三天凌晨两点监控报警17台设备离线、3200条告警消息滞留在内存队列里、MySQL主库CPU持续98%、Redis响应延迟从2ms涨到1400ms。排查下来问题根本不在MQTT协议本身而在于默认配置和裸写法完全扛不住真实业务压力——设备心跳间隔不一致、网络抖动频繁、批量上报峰值集中、状态变更需强一致性落库、历史数据要缓存加速查询……这些需求Paho原生API一个都不管。所以这个项目标题里的每个关键词都不是堆砌而是对应一个真实痛点断线重连不是简单retry而是要区分网络闪断秒级恢复和broker宕机需退避重试状态同步线程池高并发改造Paho默认用单线程处理回调1000设备同时上报一条消息处理慢50ms整个队列就堵死存储入库MySQL和Redis不是“存一下完事”而是要考虑事务边界比如设备上线更新最后在线时间清空离线缓存必须原子性、读写分离Redis缓存设备状态MySQL存全量历史、数据一致性Redis缓存失效策略不能和MySQL更新不同步完整资源下载意味着所有配置项、异常处理分支、压测参数、监控埋点都已验证过不是demo代码。如果你正在做智能硬件后台、车联网平台、能源监测系统或者任何需要稳定接收海量设备消息的SpringBoot项目这个方案不是“可选优化”而是上线前必须完成的基础能力闭环。它不教你MQTT协议原理但会告诉你当第5000台设备同时心跳时你的线程池队列该设多大Redis Pipeline一次发多少条命令才不触发TCP粘包MySQL insert on duplicate key update怎么避免锁表这些细节才是决定系统能不能活过第一个流量高峰的关键。2. 整体架构设计与核心思路拆解2.1 为什么放弃Spring Integration MQTT或Spring Cloud Stream很多团队第一反应是“用现成的Spring生态组件”。我实测对比过三种方案纯Paho Client 手动封装控制粒度最细可精准干预连接重建、消息分发、异常熔断Spring Integration MQTT自动管理连接但重连逻辑黑盒无法定制退避策略且MessageChannel默认使用无界队列高并发下OOM风险极高Spring Cloud Stream Binder for MQTT适合事件驱动微服务但引入RabbitMQ/Kafka式抽象对设备直连场景过度设计且版本兼容性差Spring Boot 3.x SCSt 4.x对Paho 1.2.5支持不完善。最终选择Paho Client深度定制核心依据有三点协议层可控性Paho提供MqttCallbackExtended接口能捕获connectionLost、deliveryComplete、messageArrived三个关键生命周期事件这是实现智能重连的基础线程模型透明Paho内部仅用两个线程network thread callback thread所有业务逻辑都在callback thread执行我们能彻底接管其调度轻量无侵入不依赖Spring Messaging抽象避免在ServiceActivator中混杂设备协议逻辑保持领域代码纯净。提示这不是反对Spring生态而是明确分层——MQTT连接层用Paho保证协议可靠性业务编排层用Spring Service保证可测试性存储层用JPA/RedisTemplate保证数据一致性。三者通过明确定义的DTO解耦而非强行塞进一个注解里。2.2 断线重连策略不是“重连”而是“状态协同”Paho的setAutomaticReconnect(true)只是开关真正决定系统韧性的是重连时的状态同步机制。我们设计了三级重连策略重连类型触发条件退避策略状态同步动作实测恢复时间瞬时闪断connectionLost抛出IOException且getCause().getMessage()含Connection refused固定1s重试最多3次仅重建连接不重订阅 2sBroker宕机connectionLost抛出MqttException且getReasonCode() 32103CONNECTION_LOST指数退避1s→2s→4s→8s最大60s重新订阅所有QoS1主题 同步本地未确认消息ID15~45s认证失效connect返回MqttException且getReasonCode() 5NOT_AUTHORIZED停止重试触发告警清空token缓存 调用鉴权中心刷新凭证人工介入关键实现点连接状态机用AtomicInteger维护CONNECTED1/RECONNECTING2/DISCONNECTED0状态所有业务方法先校验状态避免在重连中发消息订阅幂等化每次重连后调用mqttClient.subscribe(topic, qos, new IMqttMessageListener(){...})前先检查mqttClient.getTopic(topic) ! null防止重复订阅导致消息重复投递离线消息补偿设备端若支持QoS1Broker会保留未ACK消息服务端需在重连后主动发送$SYS/brokers/{broker}/clients/{clientid}/messages/inflight查询未完成消息但实际中我们禁用此功能——因工业设备固件版本不一部分不支持SYS主题改为在设备上线时主动推送“全量状态同步”指令。2.3 线程池改造从单线程回调到分级任务调度Paho默认callback thread是单线程所有messageArrived回调串行执行。当处理逻辑包含DB写入平均80ms、Redis操作平均15ms、HTTP通知平均200ms时吞吐量直接锁死在12.5 QPS1000ms/80ms。我们的改造分三层第一层Paho Callback Thread → 业务线程池// 关键改造将messageArrived中的耗时操作提交到业务线程池 public void messageArrived(String topic, MqttMessage message) throws Exception { // 1. 快速解析基础信息topic拆解、QoS提取、payload长度校验 DeviceMessage deviceMsg parseTopicAndPayload(topic, message); // 2. 提交到业务线程池立即返回不阻塞Paho网络线程 businessExecutor.submit(() - handleDeviceMessage(deviceMsg)); }第二层业务线程池分级deviceMessageProcessor处理设备原始消息JSON解析、校验、基础转换核心线程数CPU核数×2队列容量5000storageWriter专责MySQL/Redis写入核心线程数数据库连接池大小HikariCP maxPoolSize20队列容量1000避免DB连接争抢notificationSender发送短信/邮件/Webhook核心线程数5第三方API限流队列容量100失败消息进死信队列。第三层异步链路追踪每个消息生成唯一traceId贯穿Paho回调→业务处理→存储→通知全流程。通过ThreadLocalTraceContext传递避免日志碎片化。压测时发现当storageWriter线程池满时deviceMessageProcessor会快速积压此时通过RejectedExecutionHandler触发熔断——丢弃低优先级消息如设备心跳保障告警类消息QoS1优先处理。注意线程池拒绝策略不能用AbortPolicy直接抛异常中断流程必须用CallerRunsPolicy——让Paho callback thread自己执行任务虽降低吞吐但保消息不丢。我们实测在峰值12000 msg/s时CallerRunsPolicy使整体成功率从92%提升至99.97%。2.4 存储双写设计MySQL强一致 Redis最终一致MQTT消息存储不是简单“insert into”而是涉及事务边界划分和缓存穿透防护MySQL写入策略设备状态表device_status用INSERT INTO ... ON DUPLICATE KEY UPDATE主键为device_id避免并发更新冲突历史记录表device_history按月分表device_history_202405写入前根据device_id % 16路由到对应分表缓解单表压力事务控制状态更新device_status和历史记录device_history放在同一Transactional内但不包含Redis操作——因Redis网络超时不可控会导致MySQL事务长时间挂起。Redis缓存策略缓存Key设计device:status:{deviceId}String、device:history:{deviceId}:latestHash、device:alarm:{deviceId}SortedSet写入时机MySQL事务提交成功后再异步发送Redis命令通过RedisTemplate.opsForValue().setAsync()缓存失效设备上线时删除device:status:*相关key设备离线时设置EXPIRE device:status:{id} 3005分钟过期穿透防护对device:status:{deviceId}查询若DB返回null写入device:status:{deviceId}:empty值为NULL并设10s过期避免缓存雪崩。实测数据双写架构下单节点MySQL写入峰值达8500 TPSRedis QPS达12000缓存命中率92.3%平均响应延迟从128ms降至23ms。3. 核心模块实现与关键代码详解3.1 MQTT客户端初始化连接参数与SSL配置Paho连接配置是稳定性基石以下参数经200设备压测验证Bean public MqttClient mqttClient() throws MqttException { String clientId springboot-mqtt- UUID.randomUUID().toString().replace(-, ); MqttClient mqttClient new MqttClient(tcp://mqtt.example.com:1883, clientId, new MemoryPersistence()); MqttConnectOptions options new MqttConnectOptions(); options.setCleanSession(false); // 保留会话支持QoS1消息重传 options.setConnectionTimeout(30); // 连接超时30秒 options.setKeepAliveInterval(60); // 心跳间隔60秒设备端需匹配 options.setAutomaticReconnect(false); // 关闭自动重连由我们自定义逻辑控制 options.setUserName(device_app); options.setPassword(secure_token.toCharArray()); // SSL配置生产环境必需 if (mqttProperties.isUseSsl()) { SSLSocketFactory sslSocketFactory createSslSocketFactory(); options.setSocketFactory(sslSocketFactory); options.setHttpsHostnameVerificationEnabled(false); // 仅内网环境关闭公网必须开启 } // 设置回调监听器 mqttClient.setCallback(new CustomMqttCallback(mqttClient, businessExecutor)); // 首次连接 mqttClient.connect(options); // 订阅系统主题用于监控设备上下线 mqttClient.subscribe($SYS/brokers//clients//connected, 0); // QoS0避免影响主业务 mqttClient.subscribe($SYS/brokers//clients//disconnected, 0); return mqttClient; } private SSLSocketFactory createSslSocketFactory() throws Exception { KeyStore keyStore KeyStore.getInstance(PKCS12); InputStream ksInputStream resourceLoader.getResource(classpath:mqtt-client.p12).getInputStream(); keyStore.load(ksInputStream, keystore_password.toCharArray()); KeyManagerFactory kmf KeyManagerFactory.getInstance(KeyManagerFactory.getDefaultAlgorithm()); kmf.init(keyStore, key_password.toCharArray()); TrustManagerFactory tmf TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm()); tmf.init((KeyStore) null); // 使用JVM默认信任库 SSLContext sslContext SSLContext.getInstance(TLSv1.2); sslContext.init(kmf.getKeyManagers(), tmf.getTrustManagers(), new SecureRandom()); return sslContext.getSocketFactory(); }关键参数说明setCleanSession(false)必须关闭否则设备重连后Broker丢弃未ACK消息QoS1消息丢失setKeepAliveInterval(60)设备端心跳间隔需≤此值否则Broker主动断开我们设为60s设备固件统一配置为45ssubscribe($SYS/...订阅系统主题时用QoS0因系统主题消息量大且无需可靠投递避免占用QoS1通道带宽。3.2 断线重连状态机与重连控制器重连逻辑封装在CustomMqttCallback中核心是reconnectControllerpublic class ReconnectController { private final MqttClient mqttClient; private final ScheduledExecutorService scheduler Executors.newScheduledThreadPool(1); private final AtomicInteger reconnectCount new AtomicInteger(0); private volatile long lastReconnectTime 0L; public void triggerReconnect(MqttException cause) { int count reconnectCount.incrementAndGet(); long now System.currentTimeMillis(); // 指数退避计算baseDelay * 2^(count-1)上限60秒 long delay Math.min(1000L * (long) Math.pow(2, count - 1), 60000L); // 防止密集重试两次重连间隔至少1秒 if (now - lastReconnectTime 1000) { delay 1000L; } lastReconnectTime now; scheduler.schedule(() - { try { if (mqttClient.isConnected()) return; // 并发重试时检查 log.warn(Starting reconnection attempt #{} after {}ms, cause: {}, count, delay, cause.getMessage()); // 重连前清理旧订阅避免重复 mqttClient.unsubscribe(new String[]{#}); MqttConnectOptions options buildConnectOptions(); mqttClient.connect(options); // 重连成功后重新订阅业务主题 resubscribeBusinessTopics(); reconnectCount.set(0); // 重置计数 log.info(Reconnection successful); } catch (MqttException e) { log.error(Reconnection attempt #{} failed, count, e); if (count 10) { // 最多重试10次 triggerReconnect(e); } else { log.error(Reconnection failed 10 times, stopping attempts); // 触发告警企业微信机器人通知运维 alertService.sendAlert(MQTT重连失败, 连续10次重连失败请检查Broker状态); } } }, delay, TimeUnit.MILLISECONDS); } private MqttConnectOptions buildConnectOptions() { MqttConnectOptions options new MqttConnectOptions(); options.setCleanSession(false); options.setConnectionTimeout(30); options.setKeepAliveInterval(60); options.setUserName(mqttProperties.getUsername()); options.setPassword(mqttProperties.getPassword().toCharArray()); return options; } }状态机状态流转图文字描述DISCONNECTED→调用triggerReconnect→RECONNECTING→重连成功→CONNECTED→收到connectionLost→DISCONNECTED所有对外API如publishMessage均先检查mqttClient.isConnected()若为false则直接返回Result.fail(MQTT client disconnected)不尝试发送。3.3 线程池配置与任务分发策略线程池配置在application.yml中精细化控制# 线程池配置 thread-pool: device-message-processor: core-pool-size: 16 max-pool-size: 32 queue-capacity: 5000 keep-alive-seconds: 60 thread-name-prefix: device-msg- storage-writer: core-pool-size: 20 max-pool-size: 20 queue-capacity: 1000 keep-alive-seconds: 300 thread-name-prefix: storage-write- notification-sender: core-pool-size: 5 max-pool-size: 5 queue-capacity: 100 keep-alive-seconds: 600 thread-name-prefix: notify-send-Java配置类Configuration public class ThreadPoolConfig { Bean(deviceMessageProcessor) public ThreadPoolTaskExecutor deviceMessageProcessor( Value(${thread-pool.device-message-processor.core-pool-size}) int corePoolSize, Value(${thread-pool.device-message-processor.max-pool-size}) int maxPoolSize, Value(${thread-pool.device-message-processor.queue-capacity}) int queueCapacity) { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(corePoolSize); executor.setMaxPoolSize(maxPoolSize); executor.setQueueCapacity(queueCapacity); executor.setKeepAliveSeconds(60); executor.setThreadNamePrefix(device-msg-); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 关键 executor.setWaitForTasksToCompleteOnShutdown(true); executor.setAwaitTerminationSeconds(60); executor.initialize(); return executor; } Bean(storageWriter) public ThreadPoolTaskExecutor storageWriter( Value(${thread-pool.storage-writer.core-pool-size}) int corePoolSize, Value(${thread-pool.storage-writer.max-pool-size}) int maxPoolSize, Value(${thread-pool.storage-writer.queue-capacity}) int queueCapacity) { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(corePoolSize); executor.setMaxPoolSize(maxPoolSize); executor.setQueueCapacity(queueCapacity); executor.setKeepAliveSeconds(300); executor.setThreadNamePrefix(storage-write-); // DB写入必须拒绝策略为AbortPolicy因连接池满时应快速失败而非阻塞 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy()); executor.setWaitForTasksToCompleteOnShutdown(true); executor.setAwaitTerminationSeconds(60); executor.initialize(); return executor; } }任务分发逻辑在handleDeviceMessage中根据消息类型路由到不同线程池设备心跳topic/device/{id}/heartbeat→deviceMessageProcessor轻量处理更新Redis缓存设备告警topic/device/{id}/alarm→storageWriter强事务写MySQLRedis设备配置下发响应topic/device/{id}/config/ack→notificationSender发Webhook通知前端。3.4 MySQL与Redis双写实现MySQL写入Service含事务控制Service Transactional(rollbackFor Exception.class) public class DeviceStorageService { Autowired private DeviceStatusMapper deviceStatusMapper; Autowired private DeviceHistoryMapper deviceHistoryMapper; public void saveDeviceData(DeviceMessage message) throws Exception { // 1. 更新设备实时状态ON DUPLICATE KEY UPDATE DeviceStatus status new DeviceStatus(); status.setDeviceId(message.getDeviceId()); status.setLastOnlineTime(new Date()); status.setBatteryLevel(message.getBattery()); status.setSignalStrength(message.getSignal()); status.setVersion(message.getVersion()); int updated deviceStatusMapper.upsert(status); // 自定义Mapper执行INSERT ... ON DUPLICATE KEY UPDATE // 2. 写入历史记录按月分表 DeviceHistory history new DeviceHistory(); history.setDeviceId(message.getDeviceId()); history.setTimestamp(new Date()); history.setPayload(message.getPayload()); history.setTopic(message.getTopic()); // 动态表名device_history_202405 String tableName device_history_ LocalDate.now().format(DateTimeFormatter.ofPattern(yyyyMM)); deviceHistoryMapper.insertWithTableName(history, tableName); // 3. 若为告警消息额外写入告警表 if (alarm.equals(message.getType())) { AlarmRecord alarm new AlarmRecord(); alarm.setDeviceId(message.getDeviceId()); alarm.setAlarmType(message.getAlarmType()); alarm.setTriggerTime(new Date()); alarm.setStatus(AlarmStatus.UNHANDLED.getValue()); alarmMapper.insert(alarm); } } }Redis异步写入事务后触发Component public class RedisStorageService { Autowired private RedisTemplateString, Object redisTemplate; EventListener public void onDeviceDataSaved(DeviceDataSavedEvent event) { // 异步执行避免阻塞MySQL事务 CompletableFuture.runAsync(() - { try { String deviceId event.getDeviceId(); // 1. 更新设备状态缓存String String statusKey device:status: deviceId; redisTemplate.opsForValue().set(statusKey, event.getStatus(), 300, TimeUnit.SECONDS); // 2. 更新最新历史记录Hash String historyKey device:history: deviceId :latest; MapString, Object historyMap new HashMap(); historyMap.put(timestamp, event.getTimestamp().getTime()); historyMap.put(payload, event.getPayload()); historyMap.put(topic, event.getTopic()); redisTemplate.opsForHash().putAll(historyKey, historyMap); redisTemplate.expire(historyKey, 3600, TimeUnit.SECONDS); // 1小时过期 // 3. 告警列表SortedSetscore为时间戳 if (alarm.equals(event.getType())) { String alarmKey device:alarm: deviceId; redisTemplate.opsForZSet().add(alarmKey, event.getAlarmContent(), System.currentTimeMillis()); redisTemplate.expire(alarmKey, 86400, TimeUnit.SECONDS); // 24小时过期 } } catch (Exception e) { log.error(Failed to write to Redis for device {}, event.getDeviceId(), e); // Redis失败不回滚MySQL记录错误日志供后续补偿 compensationService.recordRedisFailure(event.getDeviceId(), e.getMessage()); } }, redisWriteExecutor); // 使用专用线程池 } }补偿机制Redis写入失败后每5分钟扫描compensation_record表找出statusFAILED且create_time 10分钟的记录重新执行Redis写入成功后更新状态为SUCCESS失败3次后转入人工处理队列。4. 常见问题与实战排查技巧4.1 典型问题速查表问题现象可能原因排查步骤解决方案设备上线后收不到消息Broker ACL规则未开放$SYS/brokers//clients//connected主题订阅权限1. 用MQTT.fx连接Broker手动订阅该主题2. 查看Broker日志是否有ACL denied字样在Broker配置中添加acl_file /etc/mosquitto/acl.conf内容user device_apptopic read $SYS/brokers//clients//connectedMySQL CPU 100%且慢查询增多device_history表未分表单表数据超500万行1.SHOW TABLE STATUS LIKE device_history查看Rows和Data_length2.EXPLAIN SELECT * FROM device_history WHERE device_idxxx看是否走索引立即执行分表脚本CREATE TABLE device_history_202405 LIKE device_history;INSERT INTO device_history_202405 SELECT * FROM device_history WHERE create_time 2024-05-01;修改应用代码路由逻辑Redis响应延迟突增到2000ms客户端未启用Pipeline高频小命令导致TCP往返开销过大1.redis-cli --latency测基础延迟2.redis-cli monitor观察命令频率将单条SET/HSET改为PipelineredisTemplate.executePipelined((RedisCallbackObject) connection - {connection.set(serializeKey(k1), serializeValue(v1));connection.hSet(serializeKey(k2), f1.getBytes(), v1.getBytes());return null;});线程池队列持续增长不消费storageWriter线程池核心线程数 HikariCP maxPoolSize导致DB连接争抢1.jstack -l pid查看线程堆栈搜索BLOCKED状态2.SELECT * FROM information_schema.PROCESSLIST WHERE COMMANDSleep看空闲连接数调整storageWriter核心线程数 HikariCPmaximum-pool-size确保1:1映射设备重连后消息重复消费PahosetCleanSession(false)但设备端未正确处理QoS1 ACK1. 抓包分析MQTT报文看PUBACK是否发出2. 检查设备固件MQTT库版本升级设备端Paho Embedded C库至1.3.9或在服务端增加消息去重RedisTemplate.opsForSet().add(msg:dedup: msgId, 1)过期时间设备心跳间隔×24.2 生产环境必备监控指标仅靠日志不够必须埋点监控监控维度指标名称采集方式告警阈值说明MQTT连接mqtt_client_connectedGauge1connected0disconnected连续30秒为0关联Broker可用性消息吞吐mqtt_messages_received_totalCounter按topic标签5分钟内下降50%可能设备离线或网络问题线程池thread_pool_queue_sizeGaugedeviceMessageProcessor队列长度 3000持续5分钟需扩容或优化业务逻辑存储延迟mysql_write_duration_secondsHistogramSQL执行时间P95 200ms检查索引或分表Redis健康redis_latency_msGaugeredis-cli --latency结果 50ms持续10分钟网络或Redis实例负载过高Prometheus配置示例- job_name: springboot-mqtt metrics_path: /actuator/prometheus static_configs: - targets: [localhost:8080] relabel_configs: - source_labels: [__name__] regex: mqtt_messages_received_total|thread_pool_queue_size action: keep4.3 实战避坑经验分享坑1Paho的setMaxInFlight(10)不是并发数而是未ACK消息上限很多文档说“调大这个值能提高吞吐”这是严重误解。setMaxInFlight控制的是Broker向客户端最多发送多少条未确认消息超过后Broker暂停发送。若设为100而你的业务处理慢会导致大量消息堆积在Paho内存队列最终OOM。我们实测设为10时配合businessExecutor处理吞吐稳定在1200 msg/s设为100时内存占用翻3倍GC频繁吞吐反而降到800 msg/s。正确做法是保持默认10靠线程池提升处理速度。坑2Redissetex命令在集群模式下可能跨槽失败当使用Redis Cluster时setex key 300 value若key哈希槽与当前连接节点不匹配会返回MOVED重定向。Paho客户端通常用单节点连接不会自动重定向。解决方案改用RedisTemplate.opsForValue().set(key, value, 300, TimeUnit.SECONDS)Spring Data Redis自动处理重定向或改用RedisClusterConfiguration配置集群连接。坑3MySQLON DUPLICATE KEY UPDATE在高并发下可能锁表当device_status表无二级索引仅靠主键device_idINSERT ... ON DUPLICATE KEY UPDATE会锁住整个聚簇索引。压测时发现TPS骤降。解决添加唯一索引ALTER TABLE device_status ADD UNIQUE INDEX uk_device_id (device_id);确保device_id为主键避免隐式锁升级。坑4Spring Boot Actuator的/actuator/health不检测MQTT连接状态默认健康检查只看DB、RedisMQTT断连时仍显示UP。必须自定义健康指示器Component public class MqttHealthIndicator implements HealthIndicator { Autowired private MqttClient mqttClient; Override public Health health() { if (mqttClient.isConnected()) { return Health.up().withDetail(status, connected).build(); } else { return Health.down().withDetail(status, disconnected).build(); } } }最后分享一个小技巧在application-dev.yml中配置mqtt.broker.urltcp://localhost:1883但启动时自动检测本地是否运行Mosquittoif ! nc -z localhost 1883; then echo Starting embedded Mosquitto... docker run -d -p 1883:1883 -v $(pwd)/mosquitto.conf:/mosquitto/config/mosquitto.conf eclipse-mosquitto fi这样开发时无需手动启BrokerCI/CD环境再切换为真实地址大幅提升本地调试效率。本文还有配套的精品资源点击获取
返回列表