
簡介一份面向Spring Boot開發者的MQTT集成示例工程聚焦在Java后端項目中快速接入MQTT協議實現設備與云端的數據通信。資源以可運行的輕量級demo為核心囊括Maven依賴、配置文件、消息服務封裝類以及測試代碼適合正在學習Spring Boot消息通信或物聯網數據接入的初中級開發者參考。壓縮包內共80個文件以xml配置、java源碼、class編譯文件為主另有properties配置、maven wrapper腳本及說明文檔整體僅107KB結構精簡便于快速定位關鍵代碼。內容覆蓋MQTT協議基礎、Paho客戶端集成、QoS等級說明、發布與訂閱邏輯、斷線重連及異常處理等常見環節既適合系統學習Spring Boot與MQTT的整合思路也可以作為項目初期的腳手架直接改造復用。已有1183人學習下載是理解Spring Boot整合MQTT全流程的實用參考資料。1. 你以為加個依賴就能跑 MQTTSpring Boot 的 starter 慣性反而是最大障礙很多人在 springboot 項目里第一次接觸 MQTT 時會下意識去找spring-boot-starter-mqtt我的建議是冷靜一下。Spring 生態給的是spring-integration-mqtt它包裝了 Eclipse Paho但直接用它時客戶端生命周期、斷線重連、線程池回收都是隱形的你不設置AutomaticReconnect連接一斷就再也回不來你把它塞在Bean里看日志正常一壓測就出現MqttClientPersistence鎖競爭。手里這個 demo 的解法和教科書式 demo 相反直接用 Paho 客戶端用 Spring 管配置和生命周期把發布、訂閱、重連、回調拆成獨立組件。這樣至少能少踩一半坑。適合想把 MQTT 穩定嵌進業務系統的 Java 后端開發也適合準備 springboot 面試題時弄清楚 QoS 和會話語義的人。2. MQTT 的發布/訂閱模型為什么 QoS 值得花時間摳細節2.1 發布者、訂閱者與 Broker 的三角語義MQTT 和 HTTP 的一問一答完全不同客戶端從不直接給另一個客戶端發消息一切都要經過 Broker。物聯網設備把溫度數據發到device/room1/temperature后端訂閱這個主題Broker 負責匹配主題、決定消息給誰。Topic 是分層的匹配單層#匹配多層。這里第一個坑是訂閱關系只存在于 Broker 與訂閱者之間發布者不需要知道誰在訂閱也拿不到送達回執因此業務上的“是否送達”只能靠 QoS 和業務層確認來做。很多項目出問題不是協議用錯了而是把 MQTT 當成消息隊列來設計以為消息會像 Kafka 一樣被持久化、可回溯。實際上 MQTT 的持久化能力很克制需要 Clean Session 與 Retained Message 配合而且默認行為是“會話結束即丟棄”。所以設計第一階段就要明確哪些數據是實時控制指令哪些是需要離線補發的狀態量兩類消息對 Clean Session 和 QoS 的要求完全不同。在 Spring Boot 里集成 MQTT本質上只是解決“Java 進程和 Broker 之間怎么穩定地保持一條長連接”的問題不需要創造新協議。把協議邊界先畫清楚后面代碼才不會被業務帶著亂改。2.2 QoS 0、1、2 到達語義與生產選擇QoS 是 MQTT 協議詳解里被反復講的概念但真正使用時經常出現兩種極端有人一律 QoS 2有人一律 0。判斷依據應該是“這條消息丟一條會出多大事”而不是“別人推薦我用什么”。下表是三個等級的典型行為。QoS語義是否可能重復典型場景0最多一次發出即忘會丟不重發高頻溫濕度上報、日志流1至少一次Broker 回 PUBACK有重復設備報警、平臺指令2恰好一次四次握手不重復計費、訂單、權限變更注意 QoS 是端到端語義但實現上是逐跳保證的發布端到 Broker 是一跳Broker 到訂閱端是另一跳。發布用 QoS 1、訂閱端用 QoS 0最終消息可能還是丟。所以發布和訂閱兩端的 QoS 要按鏈路里要求最高的一端統一設計不能只看 producer 配置。還有一點容易被面試題拿來考QoS 2 不會 100% 保證不重復它只是靠協議握手把重復概率壓到極低。如果業務對冪等有強要求下游消費端還是要做去重比如用消息里的 sequence 字段配合 Redis 記錄已完成消息 ID。2.3 為什么選 Eclipse Paho 而不是 Netty 自己拉協議這個項目源碼里沒有引入spring-integration-mqtt而是直接用 Paho這是有意為之。Paho Java 客戶端提供MqttClient和MqttAsyncClient兩套 API前者是阻塞式適合低頻業務封裝后者返回 token適合在 Spring 里配合回調做事件驅動。沒有特殊性能需求時不需要用 Netty 自研 MQTT 協議棧QoS、心跳、會話恢復都是極考驗邊界情況的活自己寫很難比 Paho 更穩。Paho 還有一個容易被忽略的能力MqttCallbackExtended。它比老的MqttCallback多了connectComplete(boolean reconnect, String serverURI)回調這個回調是重連成功時重新訂閱主題的正確入口。如果用MqttCallback只能在messageArrived里判斷連接狀態重連后的訂閱恢復做不干凈。所以建議一上來就實現MqttCallbackExtended把“連接完成后要做的事”集中到一個方法里這是后面 3.3 節 resubscribe 的基礎。3. Spring Boot 配置與 MqttAsyncClient 生命周期把連接管理交給容器3.1 pom.xml 只需要一個 Paho 依賴在 Spring Boot 的 Maven 工程里集成 Paho 的最簡依賴是dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version /dependency版本號按你們內部依賴管理辦法統一調整重點注意 groupId 是org.eclipse.paho不是 Spring starter。直連 Paho 而不是走 Spring Integration 模塊能保證你對連接對象有完整控制權什么時候創建、什么時候斷開、斷線后做什么都在自己代碼里不依賴框架內部 adapter 的調度邏輯。這個依賴包含MqttAsyncClient、MqttConnectOptions、MqttMessage等全部核心類但不含任何 Spring 注解。它和 Spring Boot 的關系很清晰對象生命周期交給容器業務邏輯交給 Service框架邊界不會混在一起。3.2 application.yml 配置項與配置綁定將連接參數集中在application.yml而不是散落在Value里mqtt: broker: tcp://127.0.0.1:1883 client-id: spring-boot-mqtt-demo username: admin password: changeit connect-timeout-seconds: 10 keep-alive-seconds: 60 qos: 1 topics: - device//status - system/alert對應配置類使用ConfigurationProperties綁定避免寫一長串Value(${mqtt.broker})ConfigurationProperties(prefix mqtt) public class MqttProperties { private String broker; private String clientId; private String username; private String password; private int connectTimeoutSeconds 10; private int keepAliveSeconds 60; private int qos 1; private ListString topics new ArrayList(); // getter / setter 略由 IDEA 生成 }Spring Boot 3.x 下這套綁定兼容若需要 IDE 配置提示額外引入spring-boot-configuration-processor即可。如果公司安全要求密碼不落明文可以把password用 Jasypt 加密后寫成ENC(...)密文但 Paho 拿到的仍是解碼后的明文密鑰管理要注意別和配置文件放一起。配置參數含義如下配置項含義備注brokerBroker 地址必須tcp://或ssl://開頭client-id客戶端標識同一 ID 重復連接會把前一個踢下線clean-session是否清除會話需要離線收消息時設為 falsekeep-alive-seconds心跳間隔建議大于等于 broker 要求connection-timeout建立連接超時網絡抖動環境可調大qos發布訂閱默認 QoS按業務選不是越大越好topics啟動時訂閱的主題重連后需要重新訂閱配置類上的prefix mqtt意味著client-id會自動映射到clientIdSpring 的 relaxed binding 會處理短橫線。這里不建議再用Value逐個取后期加一個連接選項就要改一遍注入的參數維護成本很高。3.3 配置類創建連接、等待連接成功、注冊擴展回調MqttConfig的核心邏輯是生成MqttAsyncClient單例并保證連接建立完成后注冊回調Configuration EnableConfigurationProperties(MqttProperties.class) public class MqttConfig { private static final Logger log LoggerFactory.getLogger(MqttConfig.class); Bean(destroyMethod disconnect) public MqttAsyncClient mqttAsyncClient(MqttProperties props) throws MqttException { MqttAsyncClient client new MqttAsyncClient(props.getBroker(), props.getClientId()); MqttConnectOptions options new MqttConnectOptions(); options.setAutomaticReconnect(true); options.setCleanSession(false); options.setConnectionTimeout(props.getConnectTimeoutSeconds()); options.setKeepAliveInterval(props.getKeepAliveSeconds()); options.setUserName(props.getUsername()); options.setPassword(props.getPassword().toCharArray()); client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { log.info(MQTT 連接完成reconnect{}, serverURI{}, reconnect, serverURI); if (reconnect) { resubscribe(client, props.getTopics(), props.getQos()); } } Override public void connectionLost(Throwable cause) { log.warn(MQTT 連接丟失{}, cause.getMessage()); } Override public void messageArrived(String topic, MqttMessage message) { log.info(收到主題 {} 的消息{}, topic, new String(message.getPayload(), StandardCharsets.UTF_8)); } Override public void deliveryComplete(IMqttDeliveryToken token) { log.debug(消息投遞完成tokenId{}, token.getMessageId()); } }); client.connect(options).waitForCompletion(15_000); resubscribe(client, props.getTopics(), props.getQos()); return client; } private void resubscribe(MqttAsyncClient client, ListString topics, int qos) { for (String topic : topics) { try { client.subscribe(topic, qos); log.info(已訂閱主題{}QoS{}, topic, qos); } catch (MqttException e) { log.error(訂閱主題 {} 失敗, topic, e); } } } }代碼邏輯說明Bean(destroyMethod disconnect)讓 Spring 在關閉應用時自動調用 Paho 的disconnect()setAutomaticReconnect(true)打開 Paho 內置重連讓它負責斷網后的 TCP 重連setCleanSession(false)讓 Broker 保留離線期間發給該 client 的消息前提是連接前已經建立持久會話。connectComplete和啟動后的resubscribe調用了同一方法確保斷線重連后訂閱不丟。一個容易忽略的細節client.connect(options)返回IMqttToken這里調用waitForCompletion(15_000)是為了讓應用啟動階段就知道 Broker 是否可達。如果要求 Broker 暫時不可用時應用也能起來就不要同步等待改成在回調里判斷連接失敗但這樣下游注入 client 的 Bean 可能拿到未連接對象所以 demo 選擇同步等待。硬要兼容啟動時 Broker 未就緒的場景可以加一層SmartLifecycle推遲發布這類做法對中小項目往往過度設計不用照搬。4. 發布與訂閱的工程化封裝不要直接在 Controller 里 new MqttMessage4.1 用 MqttPublisher 統一處理發布參數與異常消息發布經常被順手寫在業務方法里等到需要統一加 retained 標志、異步確認、埋點日志時就后悔了。這里抽出一個MqttPublisher對外只暴露明確參數Service public class MqttPublisher { private static final Logger log LoggerFactory.getLogger(MqttPublisher.class); private final MqttAsyncClient client; public MqttPublisher(MqttAsyncClient client) { this.client client; } public PublishResult publish(String topic, String payload, int qos, boolean retained) { MqttMessage message new MqttMessage(payload.getBytes(StandardCharsets.UTF_8)); message.setQos(qos); message.setRetained(retained); try { IMqttDeliveryToken token client.publish(topic, message); token.waitForCompletion(3_000); log.info(發布成功 topic{}, qos{}, retained{}, topic, qos, retained); return new PublishResult(true, token.getMessageId(), null); } catch (MqttException e) { log.error(發布失敗 topic{}, qos{}, topic, qos, e); return new PublishResult(false, -1, e.getMessage()); } } }PublishResult是一個簡單 POJO保存成功標志、消息 ID 和錯誤信息。這里的waitForCompletion(3_000)把異步發布變成最多阻塞 3 秒的同步調用適合業務方法需要立即確認協議層是否入隊的場景如果只需發出去不管結果可以去掉這一行直接返回 token性能更好但異常感知滯后。retained參數容易被忽略。Retained 表示 Broker 要保留這條消息作為“最新值”當新訂閱者訂閱該主題時立即推送。拿它做設備狀態快照很合適但如果把一次性控制指令標記為 retained會帶來詭異的重放問題新設備上線訂閱主題時會收到一條早已過期的指令。所以每次調用 publish 都顯式傳 retained默認false最安全。4.2 消息進入 Spring 領域回調里只做分發Paho 的messageArrived是在客戶端內部線程上調用的如果在回調里直接寫數據庫、調外部接口一旦阻塞后續消息全部排隊心跳也可能超時。常見做法是把“收消息”和“處理消息”拆開。回調只負責包裝成 Spring 事件業務用EventListener異步消費。MqttConfig中的 callback 改成事件發布Override public void messageArrived(String topic, MqttMessage message) { eventPublisher.publishEvent(new MqttMessageEvent(this, topic, message.getPayload(), message.getQos())); }MqttMessageEvent定義成簡單事件對象public class MqttMessageEvent extends ApplicationEvent { private final String topic; private final byte[] payload; private final int qos; public MqttMessageEvent(Object source, String topic, byte[] payload, int qos) { super(source); this.topic topic; this.payload payload; this.qos qos; } // getter 略 }再寫一個監聽組件按主題前綴分流Component public class MqttEventHandler { private static final Logger log LoggerFactory.getLogger(MqttEventHandler.class); Async(mqttTaskExecutor) EventListener public void onMessage(MqttMessageEvent event) { String topic event.getTopic(); String payload new String(event.getPayload(), StandardCharsets.UTF_8); if (topic.startsWith(device/)) { handleDeviceData(topic, payload); } else if (topic.equals(system/alert)) { handleAlert(payload); } } }這里用Async轉發到獨立線程池避免阻塞 Paho 回調線程。線程池 bean 可以用ThreadPoolTaskExecutor定義核心線程數不要超過連接數太多避免消息風暴時線程爆炸。注意 Spring 事件默認同步Async要生效需要顯式加EnableAsync。如果不引入事件機制也可以直接注入一個ExecutorService在messageArrived里提交任務。兩種方式等價但事件不會讓 MQTT 層反向依賴業務包耦合度更低。4.3 Clean Sessionfalse 與重連后的重復訂閱很多 MQTT 問題出現在重啟上客戶端第一次訂閱成功后Broker 記錄了會話應用重啟時用同一個 clientId 連接但設置了cleanSessiontrue或者沒有恢復訂閱消息就收不到了。MqttCallbackExtended把重連路徑單獨留了出來在connectComplete里判斷reconnect參數只在重連成功后重新執行訂閱。MqttAsyncClient.subscribe在斷線狀態下調用不會立刻拋異常它會先入隊等連接恢復后自動發送。所以有些項目不重訂閱也能恢復這依賴于automaticReconnect的隊列行為。問題是如果 Broker 側會話因超時被清理入隊的訂閱請求執行時可能拿到新 session 的初始狀態最終不可預期。顯式重訂閱能讓訂閱關系與應用配置保持一致日志里也能看到每次重連后的訂閱記錄排查問題更直觀。重連重訂閱還有一個順序問題connectComplete觸發時Broker 已經接受連接但未必完成此前訂閱的恢復。直接調用subscribePaho 會在連接完全就緒后發送 SUBSCRIBE這個順序是對的。不要在connectionLost里馬上重連那會打亂 Paho 自己的重連狀態機。另一個常見坑是同一個MqttAsyncClient并發發布大量消息時Paho 內部有 token 隊列若客戶端已被disconnect()后繼續 publish會拋MqttException。所以發布端必須捕獲異常并做降級不能把 MqttException 直接拋給 Controller 層。上面PublishResult就是簡單兜底。5. 用 10 分鐘驗證整套連接斷網重連與訂閱恢復5.1 本地起一個 Broker先用 Docker 起一個輕量 broker 做驗證Mosquitto 足夠docker run -d --name mqtt-demo -p 1883:1883 eclipse-mosquitto:2如果容器默認配置拒絕匿名登錄掛載一個最小配置cat /tmp/mosquitto.conf EOF listener 1883 allow_anonymous true EOF docker run -d --name mqtt-demo -p 1883:1883 \ -v /tmp/mosquitto.conf:/mosquitto/config/mosquitto.conf \ eclipse-mosquitto:2啟動 Spring Boot 應用后日志里應該看到“已訂閱主題device//status”。然后用mosquitto_pub發一條測試消息mosquitto_pub -h 127.0.0.1 -p 1883 -t device/room1/status -m {online:true} -q 1如果應用日志出現對應的“收到主題”記錄說明訂閱和回調鏈路是通的。這里用真實 broker 驗證比單測更可靠能覆蓋 Paho 客戶端與服務端握手、心跳、PUBACK 全套路徑。5.2 模擬斷線確認自動重連和重訂閱執行docker stop mqtt-demo日志應出現“MQTT 連接丟失”再執行docker start mqtt-demo日志出現MQTT 連接完成reconnecttrue且再次輸出“已訂閱主題”。這一步驗證MqttCallbackExtended的重連路徑也是生產環境最容易出問題的一段。如果只看到連接完成而沒重新訂閱檢查topics是否為空以及connectComplete里是否真的調用了resubscribe。最后從 demo.zip 的目錄結構說一個落地建議src/main/java下的 config、service、domain 分層不要打散。MqttConfig只管連接與回調MqttPublisher只管發布MqttMessageEvent只管傳遞消息設備業務邏輯全部放在EventListener之后。這樣后面加一個device//config主題只需要改配置文件并新增一個 handler 方法不用動連接層。部署上線前記得把 yml 里的mqtt.broker換成生產 broker 地址并確認 ACL 允許當前用戶名發布和訂閱對應主題否則即使連接成功也會在 PUBACK 階段被 broker 以權限問題斷開。本文還有配套的精品資源點擊獲取