spring boot 连接emqx并实现发布订阅
1、安装依赖dependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId version7.1.0/version scopecompile/scope /dependency !-- Source: https://mvnrepository.com/artifact/org.eclipse.paho/org.eclipse.paho.client.mqttv3 -- dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version scopecompile/scope /dependency2、配置连接信息#MQTT???? #MQTT-??? spring.mqtt.usernamexxxxxx #MQTT-?? spring.mqtt.passwordxxxxx #MQTT-??????????????????????tcp://127.0.0.1:61613?tcp://47.123.33.66:61613 spring.mqtt.urltcp://127.0.0.1:1883 #MQTT-??????????ID spring.mqtt.client-idemq2 #MQTT-????????????????????? spring.mqtt.topicxiaomingming #timeout ?????? spring.mqtt.timeout20 #keep alive spring.mqtt.keep-alive20 spring.mqtt.qos23、新建mqtt配置类package com.example.springbootemq.config; import org.eclipse.paho.client.mqttv3.MqttClient; import org.eclipse.paho.client.mqttv3.MqttConnectOptions; import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Configuration; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; Component Configuration public class MqttConfig { Value(${spring.mqtt.username}) private String username; Value(${spring.mqtt.password}) private String password; Value(${spring.mqtt.url}) private String hostUrl; Value(${spring.mqtt.client-id}) private String clientId; Value(${spring.mqtt.timeout}) private Integer timeout; Value(${spring.mqtt.keep-alive}) private Integer keepAlive; Value(${spring.mqtt.qos}) private Integer qos; Value(${spring.mqtt.topic}) private String topic; private MqttClient client; /** * 项目启动时自动连接 MQTT */ PostConstruct public void init() { connect(); } //断线手动重连重新订阅 public void reConnectSubscribe() { while (true){ try { log.warn(开始重连重订阅); Thread.sleep(2000); this.client.connect(connOpts); //订阅 this.client.subscribe(myTest,2); break; }catch (Exception ex){ ex.printStackTrace(); } } } /** * 连接 MQTT */ public void connect() { try { client new MqttClient(hostUrl, clientId, new MemoryPersistence()); // MQTT 连接选项 MqttConnectOptions connOpts new MqttConnectOptions(); connOpts.setUserName(username); connOpts.setPassword(password.toCharArray()); // 保留会话 connOpts.setCleanSession(true); // 设置超时时间单位秒 connOpts.setConnectionTimeout(timeout); // 设置心跳时间单位秒表示服务器每隔1.5*20秒的时间向客户端发送心跳判断客户端是否在线 connOpts.setKeepAliveInterval(keepAlive); // 设置回调 client.setCallback(new OnMessageCallback()); // 建立连接 client.connect(connOpts); //订阅 client.subscribe(topic,2); } catch (MqttException me) { System.out.println(reason me.getReasonCode()); System.out.println(msg me.getMessage()); System.out.println(loc me.getLocalizedMessage()); System.out.println(cause me.getCause()); System.out.println(excep me); me.printStackTrace(); } } /** * 订阅 * * param topic 主题 */ public void subscribe(String topic) { try { client.subscribe(topic, qos); } catch (MqttException me) { me.printStackTrace(); } } /** * 消息发布 * * param topic 主题 * param data 消息 */ public void publish(String topic, String data) { try { MqttMessage message new MqttMessage(data.getBytes()); message.setQos(qos); // 消息服务质量等级 message.setRetained(true); // 保留消息 client.publish(topic, message); } catch (MqttException me) { me.printStackTrace(); } } /** * 断开连接 */ public void disconnect() { try { client.disconnect(); client.close(); } catch (MqttException me) { me.printStackTrace(); } } }4、定义消息回调类OnMessageCallbackpackage com.example.springbootemq.config; import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken; import org.eclipse.paho.client.mqttv3.MqttCallback; import org.eclipse.paho.client.mqttv3.MqttMessage; public class OnMessageCallback implements MqttCallback { private MqttConfig mqttConfig; //构造函数注入对象 public OnMessageCallback(MqttConfig mqttConfig) { this.mqttConfig mqttConfig; } Override public void connectionLost(Throwable cause) { // 连接丢失后一般在这里面进行重连 System.out.println(连接断开可以做重连); this.mqttConfig.reConnectSubscribe(); } Override public void messageArrived(String topic, MqttMessage message) throws Exception { // subscribe后得到的消息会执行到这里面 System.out.println(接收消息主题: topic); System.out.println(接收消息Qos: message.getQos()); System.out.println(接收消息内容: new String(message.getPayload())); } Override public void deliveryComplete(IMqttDeliveryToken token) { System.out.println(deliveryComplete--------- token.isComplete()); } }5、控制器测试package com.example.springbootemq.controller; import com.example.springbootemq.config.MqttConfig; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; RestController RequestMapping(/mqtt) public class MqttController { Resource private MqttConfig mqttConfig; GetMapping(/send/{message}) public void sendMessage(PathVariable(message) String message) { String subTopic testtopic/#; String pubTopic testtopic/1; String topicTestzmtest; String data hello MQTT test; mqttConfig.publish(topicTest, message); // // 订阅 // mqttConfig.subscribe(subTopic); // // 发布消息 // mqttConfig.publish(pubTopic, data); // 断开连接 // mqttConfig.disconnect(); } GetMapping(/receive) public void receiveMessage() { String subTopic testtopic/#; String pubTopic testtopic/1; String topicTestzmtest; String data hello MQTT test; // mqttConfig.publish(topicTest, data); // // 订阅 mqttConfig.subscribe(topicTest); // // 发布消息 // mqttConfig.publish(pubTopic, data); // 断开连接 // mqttConfig.disconnect(); } }6、注意下面是高级版的重连和重新订阅1设置自动重连//是否自动重连 connOpts.setAutomaticReconnect(true);2不用新建回调类直接写client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { // 连接/重连成功后订阅 System.out.println(serverURI); try { client.subscribe(test1,qos); client.subscribe(test2,qos); } catch (Exception e) { e.printStackTrace(); } } // 连接丢失后一般在这里面进行重连 Override public void connectionLost(Throwable throwable) { System.out.println(连接丢失1); } Override public void messageArrived(String topic, MqttMessage message) { try { System.out.println(接收消息内容: new String(message.getPayload())); }catch (Exception ex){ ex.printStackTrace(); } } Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println(111111); } });3完整MqttConfig配置类package com.example.superior_conjuncte_iot_rtspvideo.system_config.config; import com.example.superior_conjuncte_iot_rtspvideo.system_config.handle.OnMessageCallback; import jakarta.annotation.PostConstruct; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.eclipse.paho.client.mqttv3.*; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import java.util.Properties; Component Configuration Slf4j public class MqttConfig { Resource Lazy private OnMessageCallback onMessageCallback; Value(${spring.mqtt.username}) private String username; Value(${spring.mqtt.password}) private String password; Value(${spring.mqtt.url}) private String hostUrl; Value(${spring.mqtt.client-id}) private String clientId; Value(${spring.mqtt.timeout}) private Integer timeout; Value(${spring.mqtt.keep-alive}) private Integer keepAlive; Value(${spring.mqtt.qos}) private Integer qos; Value(${kafka.kafkaTopic}) private String kafkaTopic; private String topiccamera/1/video; Resource private KafkaAttributeConfig kafkaAttributeConfig; private KafkaProducerString, String kafkaVideoProducernull; private MqttClient client; /** * 项目启动时自动连接 MQTT */ PostConstruct public void init() { //初始化kafka initKafkaProducer(); connect(); } public void initKafkaProducer(){ Properties propskafkaAttributeConfig.initProductConfig(); this.kafkaVideoProducer new KafkaProducer(props); } //发送到kafka生产者 public void sendKafkaVideo(String key,String message){ ProducerRecordString,String record new ProducerRecord(kafkaTopic,key,message); kafkaVideoProducer.send(record); } //断线手动重连重新订阅 public void reConnectSubscribe() { while (true){ try { log.warn(开始重连重订阅); Thread.sleep(2000); this.connect(); // //订阅 // this.client.subscribe(myTest,2); break; }catch (Exception ex){ ex.printStackTrace(); } } } /** * 连接 MQTT */ public void connect() { try { client new MqttClient(hostUrl, clientId, new MemoryPersistence()); // MQTT 连接选项 MqttConnectOptions connOpts new MqttConnectOptions(); connOpts.setUserName(username); connOpts.setPassword(password.toCharArray()); // 保留会话 connOpts.setCleanSession(true); // 设置超时时间单位秒 connOpts.setConnectionTimeout(timeout); // 设置心跳时间单位秒表示服务器每隔1.5*20秒的时间向客户端发送心跳判断客户端是否在线 connOpts.setKeepAliveInterval(keepAlive); //是否自动重连 connOpts.setAutomaticReconnect(true); // 设置回调 // client.setCallback(onMessageCallback); client.setCallback(new MqttCallbackExtended() { Override public void connectComplete(boolean reconnect, String serverURI) { // 连接/重连成功后订阅 System.out.println(serverURI); try { client.subscribe(test1,qos); client.subscribe(test2,qos); } catch (Exception e) { e.printStackTrace(); } } // 连接丢失后一般在这里面进行重连 Override public void connectionLost(Throwable throwable) { System.out.println(连接丢失1); } Override public void messageArrived(String topic, MqttMessage message) { try { System.out.println(接收消息内容: new String(message.getPayload())); }catch (Exception ex){ ex.printStackTrace(); } } Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println(111111); } }); // 建立连接 client.connect(connOpts); //订阅 // client.subscribe(topic,qos); // client.subscribe(test1,qos); // client.subscribe(test2,qos); } catch (MqttException me) { log.error(MQTT连接失败{}, me.getMessage()); } } /** * 断开连接 */ public void disconnect() { try { client.disconnect(); client.close(); } catch (MqttException ex) { log.error(MQTT断开连接失败{}, ex.getMessage()); } } }

相关新闻

运维手动改了一行数据,Seata 回滚时把它覆盖了:AT 模式全局锁和 undo_log 的 3 个被忽略的边界

运维手动改了一行数据,Seata 回滚时把它覆盖了:AT 模式全局锁和 undo_log 的 3 个被忽略的边界

title: 运维手动改了一行数据,Seata 回滚时把它覆盖了:AT 模式全局锁和 undo_log 的 3 个被忽略的边界 date: 2026-08-22 category: 分布式事务 tags: [Java, 后端, 分布式事务, Seata, 微服务]我们订单和库存两个服务用 Seata AT 模式做分布式事务。AT …

2026/8/22 16:15:32 阅读更多 →
看门狗每 10 秒续一次期,业务却跑了 40 秒:Redisson 锁续期源码里藏着的 2 个默认陷阱

看门狗每 10 秒续一次期,业务却跑了 40 秒:Redisson 锁续期源码里藏着的 2 个默认陷阱

title: 看门狗每 10 秒续一次期,业务却跑了 40 秒:Redisson 锁续期源码里藏着的 2 个默认陷阱 date: 2026-08-22 category: 分布式锁 tags: [Java, 后端, Redis, Redisson, 并发编程]我们订单服务用 Redisson 做分布式锁,锁一段「扣减库存 生…

2026/8/22 16:15:32 阅读更多 →
设计模式结构型——桥接模式

设计模式结构型——桥接模式

目录 什么是桥接模式 桥接模式的实现 桥接模式角色 桥接模式类图 桥接模式举例 桥接模式代码实现 桥接模式的特点 优点 缺点 使用场景 注意事项 什么是桥接模式 桥接(Bridge)模式是用于把抽象化与实现化解耦,使得二者可以独立变化…

2026/8/22 16:14:32 阅读更多 →

最新新闻

一个免费油猴脚本,把 NGA 论坛的冗余信息一键清掉(附快捷键与安装步骤)

一个免费油猴脚本,把 NGA 论坛的冗余信息一键清掉(附快捷键与安装步骤)

一个免费油猴脚本,把 NGA 论坛的冗余信息一键清掉(附快捷键与安装步骤) 【免费下载链接】NGA-BBS-Script NGA论坛增强脚本,给你完全不一样的浏览体验 项目地址: https://gitcode.com/gh_mirrors/ng/NGA-BBS-Script 在 NGA …

2026/8/22 16:59:50 阅读更多 →
TradingAgents 无 GPU 部署实战:三步让 LLM 多智能体团队跑通 AAPL 回测

TradingAgents 无 GPU 部署实战:三步让 LLM 多智能体团队跑通 AAPL 回测

TradingAgents 无 GPU 部署实战:三步让 LLM 多智能体团队跑通 AAPL 回测 【免费下载链接】TradingAgents-AI.github.io TradingAgents: Multi-Agents LLM Financial Trading Framework 项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-AI.github…

2026/8/22 16:59:50 阅读更多 →
从TRAXX机车到BSP 25T:模块化架构与工程复用的工业软件实践

从TRAXX机车到BSP 25T:模块化架构与工程复用的工业软件实践

如果你在铁路技术论坛或开发者社区看到“BSP原型车”这个词,可能会有点困惑——这听起来像是一个软件项目或硬件原型。但今天我们要聊的,是一个在铁路工业软件、仿真建模和数字孪生领域极具代表性的经典案例:如何通过一个真实的机车车型&…

2026/8/22 16:59:50 阅读更多 →
数学建模在数字版权保护中的应用:从指纹算法到动态决策系统

数学建模在数字版权保护中的应用:从指纹算法到动态决策系统

1. 项目概述:当数学建模遇上数字版权保护最近刚带学生打完“深圳杯”数学建模挑战赛,B题“电子资源版权保护问题”给我留下了挺深的印象。这题目出得相当有水平,它没有停留在传统的版权法理探讨上,而是直接把一个复杂的现实问题&a…

2026/8/22 16:59:50 阅读更多 →
免训练AI模型Nori:表格数据缺失值填充与预测实战指南

免训练AI模型Nori:表格数据缺失值填充与预测实战指南

这次我们来看一个专门处理结构化数据预测的开源项目——Synthefy 推出的 Nori 模型。它不是生成图片或语音的 AI,而是直接帮你预测表格、数据库里那些“缺失值”或“未来值”的工具。简单说,给你一张销售表,它能预测下个月的销售额&#xff1…

2026/8/22 16:59:50 阅读更多 →
用 MediaInfo 把视频元数据拆得明明白白

用 MediaInfo 把视频元数据拆得明明白白

用 MediaInfo 把视频元数据拆得明明白白 【免费下载链接】MediaInfo Convenient unified display of the most relevant technical and tag data for video and audio files. 项目地址: https://gitcode.com/gh_mirrors/me/MediaInfo 接到一个来路不明的视频文件&#x…

2026/8/22 16:58:49 阅读更多 →

日新闻

沉金PCB工艺实战指南:从设计到SMT焊接的可靠性保障

沉金PCB工艺实战指南:从设计到SMT焊接的可靠性保障

在电子硬件开发领域,PCB(印制电路板)的沉金工艺是提升产品可靠性和焊接质量的关键环节。对于需要高密度互连、长期稳定运行或高频信号传输的板卡,如“黍姐仿通行证”这类可能涉及身份识别、数据交互的硬件项目,选择正确…

2026/8/22 0:00:11 阅读更多 →
电气考研电路八月强化四步法:从知识体系到真题实战的闭环攻略

电气考研电路八月强化四步法:从知识体系到真题实战的闭环攻略

这次我们来看一个针对电气考研电路科目的学习规划项目。它不是软件工具,而是一套聚焦于8月份关键节点的备考策略。对于电气工程考研的同学来说,电路分析是专业课的重中之重,也是拉开分差的关键。进入8月,复习进入强化阶段&#xf…

2026/8/22 0:00:11 阅读更多 →
消除AI代码的“AI味”:Claude Code设计优化技能配置与实战指南

消除AI代码的“AI味”:Claude Code设计优化技能配置与实战指南

大家好,我是专注于前端开发与AI工具实践的技术博主。在日常使用 Claude Code 等AI编程助手时,你是否也遇到过这样的困扰:生成的代码功能上没问题,但代码风格、组件设计、交互逻辑总透着一股“AI味”——布局单调、样式简陋、交互生…

2026/8/22 0:00:11 阅读更多 →

周新闻

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

如果你是一名开发者,最近可能已经感受到了AI大模型正在从“玩具”变成“生产力工具”的强烈信号。从代码补全到智能Agent,从本地部署到云端API,我们正处在一个技术栈快速重构的节点。然而,面对层出不穷的模型、框架和工具&#xf…

2026/8/21 3:21:33 阅读更多 →
工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

第四篇:反射——高频能量撞墙之后会发生什么? —— 你以为信号已经过去了,其实它正在回来打你 老Q的现场笔记 第五季,我们正式进入工业神经系统层。这里不再是单个设备的战斗,而是整个工厂“经脉”层面的秩序之战。从这一篇开始,你将第一次看清:看似简单的信号传播,背…

2026/8/22 8:09:09 阅读更多 →
【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

2026/8/21 6:07:56 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/21 16:42:28 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/22 7:31:03 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/22 3:22:48 阅读更多 →