-
MQTT协议的核心设计理念是什么?
-
前言 随着物联网的火热,MQTT的应用逐渐增多 曾经也有幸使用过mqtt,今天正好总结下MQTT的使用; 一、MQTT是什么? 可以把他理解为,也是一种mq消息,设计简单且轻量级,通讯报文开销小,占用的网络带宽和资源较少,适用于低带宽、不稳定网络环境下的通讯。 MQTT采用发布/订阅模式,分为发布者和订阅者两个角色,需要一个中介来协调发布者和订阅者之间的消息传递,这个中介就是MQTT代理(Broker)。 MQTT协议在物联网领域应用广泛,包括智能家居、工业自动化、智能交通系统等。 个人简单总结: 每个客户端可以订阅一个或者多个主题(发消息,收消息) 每个客户端不订阅主题,也可以发送主题消息(只接受消息,不发送消息) 客户端A发送消息给客户端B流程为: 客户端A>>>Broker>>>客户端B --- 前置条件: a: 客户端A 发送主题消息,且与客户端B的订阅主题一致 b: 客户端B 订阅主题 二、继承步骤 1.安装MQTT 这里直接采用windows版本,解压版,比较快 下载地址 MQTT-windows版本 解压后,在bin文件下执行运行命令 .\emqx console 访问MQTT管理页面 http://localhost:18083/#/ 用户名密码admin/public 2.创建项目,引入依赖 大致分为如下步骤: yml配置 主题 用户名 密码 根据配置创建客户端实例,实例订阅主题 实现 MqttCallback 接口 1. 重连处理 connectionLost 2. 消息接受处理 messageArrived 3. 消息发生成功处理 deliveryComplete 根据客户端信息发送某个主题的消息 3. 对应步骤2的代码 yml配置 server: port: 8081 # 下面这里要看你自己的需求 customer: mqtt: broker: tcp://127.0.0.1:1883 clientList: #发布客户端ID - clientId: nxys_service #监听主题 同时订阅多个主题使用 - 分割开 subscribeTopic: mqtt/publish #用户名 userName: admin #密码 password: public #接受客户端ID - clientId: receive_service #监听主题 同时订阅多个主题使用 - 分割开 subscribeTopic: mqtt/receive #用户名 userName: admin #密码 password: public 实例信息获取 /** * Mqtt配置类 */ @Data @Configuration @ConfigurationProperties(prefix = "customer.mqtt") public class MqttConfig { /** * mqtt broker地址 */ String broker; /** * 需要创建的MQTT客户端 */ List<MqttClient> clientList; } /** * MQTT客户端 */ @Data public class MqttClient { /** * 客户端ID */ private String clientId; /** * 监听主题 */ private String subscribeTopic; /** * 用户名 */ private String userName; /** * 密码 */ private String password; } 根据信息创建实例,订阅主题 /** * MQTT客户端创建 */ @Component @Slf4j public class MqttClientCreate { @Resource private MqttClientManager mqttClientManager; @Autowired private MqttConfig mqttConfig; /** * 创建MQTT客户端 */ @PostConstruct public void createMqttClient() { List<MqttClient> mqttClientList = mqttConfig.getClientList(); for (MqttClient mqttClient : mqttClientList) { log.info("{}", mqttClient); //创建客户端,客户端ID:demo,回调类跟客户端ID一致 mqttClientManager.createMqttClient(mqttClient.getClientId(), mqttClient.getSubscribeTopic(), mqttClient.getUserName(), mqttClient.getPassword()); } } } /** * MQTT客户端管理类,如果客户端非常多后续可入redis缓存 */ @Slf4j @Component public class MqttClientManager { @Value("${customer.mqtt.broker}") private String mqttBroker; @Resource private MqttCallBackContext mqttCallBackContext; /** * 存储MQTT客户端 */ public static Map<String, MqttClient> MQTT_CLIENT_MAP = new ConcurrentHashMap<>(); public static MqttClient getMqttClientById(String clientId) { return MQTT_CLIENT_MAP.get(clientId); } /** * 创建mqtt客户端 * * @param clientId 客户端ID * @param subscribeTopic 订阅主题,可为空 * @param userName 用户名,可为空 * @param password 密码,可为空 * @return mqtt客户端 */ public void createMqttClient(String clientId, String subscribeTopic, String userName, String password) { MemoryPersistence persistence = new MemoryPersistence(); try { MqttClient client = new MqttClient(mqttBroker, clientId, persistence); MqttConnectOptions connOpts = new MqttConnectOptions(); if (null != userName && !"".equals(userName)) { connOpts.setUserName(userName); } if (null != password && !"".equals(password)) { connOpts.setPassword(password.toCharArray()); } connOpts.setCleanSession(true); if (null != subscribeTopic && !"".equals(subscribeTopic)) { AbsMqttCallBack callBack = mqttCallBackContext.getCallBack(clientId); if (null == callBack) { callBack = mqttCallBackContext.getCallBack("default"); } callBack.setClientId(clientId); callBack.setConnectOptions(connOpts); client.setCallback(callBack); } //连接mqtt服务端broker client.connect(connOpts); // 订阅主题 if (null != subscribeTopic && !"".equals(subscribeTopic)) { if (subscribeTopic.contains("-")) client.subscribe(subscribeTopic.split("-")); else // if (!subscribeTopic.equals("mqtt/receive")) { client.subscribe(subscribeTopic); } } MQTT_CLIENT_MAP.putIfAbsent(clientId, client); } catch (MqttException e) { log.error("Create mqttClient failed!", e); } } } 实现 MqttCallback 接口 /** * MQTT回调抽象类 */ @Slf4j public abstract class AbsMqttCallBack implements MqttCallback { private String clientId; private MqttConnectOptions connectOptions; public String getClientId() { return clientId; } public void setClientId(String clientId) { this.clientId = clientId; } public MqttConnectOptions getConnectOptions() { return connectOptions; } public void setConnectOptions(MqttConnectOptions connectOptions) { this.connectOptions = connectOptions; } /** * 失去连接操作,进行重连 * * @param throwable 异常 */ @Override public void connectionLost(Throwable throwable) { try { if (null != clientId) { if (null != dconnectOptions) { MqttClientManager.getMqttClientById(clientId).connect(connectOptions); } else { MqttClientManager.getMqttClientById(clientId).connect(); } } } catch (Exception e) { log.error("{} reconnect failed!", e); } } /** * 接收订阅消息 * @param topic 主题 * @param mqttMessage 接收消息 * @throws Exception 异常 */ @Override public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception { String content = new String(mqttMessage.getPayload()); handleReceiveMessage(topic, content); } /** * 消息发送成功 * * @param iMqttDeliveryToken toke */ @Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { log.info("消息发送成功"); } /** * 处理接收的消息 * @param topic 主题 * @param message 消息内容 */ protected abstract void handleReceiveMessage(String topic, String message); } /** * 默认回调 */ @Slf4j @Component("default") public class DefaultMqttCallBack extends AbsMqttCallBack { /** * @param topic 主题 * @param message 消息内容 */ @Override protected void handleReceiveMessage(String topic, String message) { log.info("接收到主题---{}", topic); log.info("接收到消息---{}", message); // 你自己的消息处理业务 } } /** * MQTT订阅回调环境类 */ @Component @Slf4j public class MqttCallBackContext { private final Map<String, AbsMqttCallBack> callBackMap = new ConcurrentHashMap<>(); /** * 默认构造函数 * * @param callBackMap 回调集合 */ public MqttCallBackContext(Map<String, AbsMqttCallBack> callBackMap) { this.callBackMap.clear(); this.callBackMap.putAll(callBackMap); } /** * 获取MQTT回调类 * * @param clientId 客户端ID * @return MQTT回调类 */ public AbsMqttCallBack getCallBack(String clientId) { return this.callBackMap.get(clientId); } } 发送消息 @RestController public class SendController { @Resource MqttClientManager mqttClientManager; @RequestMapping("/sendMessage") public String sendMessage(String topic){ try { MqttMessage mqttMessage = new MqttMessage("你好".getBytes()); mqttClientManager.getMqttClientById("nxys_service").publish(topic,mqttMessage); return "发送成功"; } catch (Exception e) { e.printStackTrace(); return "发送失败"; } } } 3 测试 启动订阅,查看MQTT 管理页面 测试发送消息,查看发送情况,接受情况 http://localhost:8081/sendMessage?topic=mqtt/receive 总结 文中涉及的所有代码: MQTT-Demo mqtt 启动后访问地址 http://localhost:18083/#/ 用户名/密码: admin/public 每个客户端可以订阅一个或者多个主题 每个客户端不订阅主题,也可以发送主题消息 客户端A发送消息给客户端B流程为: 客户端A>>>Broker>>>客户端B --- 前置条件: a: 客户端A 发送主题消息,且与客户端B的订阅主题一致 b: 客户端B 订阅主题 1 2 3 4 5 mqtt启动命令 在bin目录下,cmd 执行 .\emqx console ———————————————— 版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。 原文链接:https://blog.csdn.net/qq_32419139/article/details/136176948
-
PC根据文档可以下载根据文档这里格式不知道有没有错,下面第一个红色框是请求红色是发送的,绿色是接收,不明白400错误那里怎么解决,我用的ESP8266透传,已经接上了https 443接口
yd_247833416
发表于2024-03-11 23:57:58
2024-03-11 23:57:58
最后回复
yd_237139617
2025-03-17 09:32:19
292 10 -
MQTT消息的存储和检索既涉及到消息的持久化方法,也包括后续如何高效地检索这些消息。1、选择合适的存储介质,利用关系型数据库如MySQL、PostgreSQL可实现消息数据的结构化存储和事务性管理。关系型数据库适合事务性强、数据结构复杂的场景,但可能需针对并发和性能进行优化措施;2、设计消息存储架构,存储层需要考虑分布式系统的设计,以保持高可用性和扩展性。借助消息队列如Kafka可以缓存大量的实时数据,并提供持久化。消息存储系统应有策略处理消息的重复、丢失、序列化等问题,并通过负载均衡、分区等技术提高系统吞吐量。数据模型设计应与MQTT协议的主题和QoS等级紧密结合,并适应不同大小和格式的消息。数据备份、副本和恢复策略也需要规划以防止数据丢失;3、实现高效的检索机制,高效检索机制涉及索引策略和查询优化技巧。 在这些步骤中,设计消息存储架构至关重要,它需要确保数据既持久化又能应对高并发的读写需求。4、整合并测试系统,系统整合将存储和检索机制与MQTT代理相连,重点在于确保消息流转无缝且高效。经历全面的压力测试和性能调优之后,可以确信系统能够在实际运行环境中稳定工作。测试阶段也需验证数据一致性和备份恢复流程的有效性。
-
什么是 MQTT协议?MQTT 全称(Message Queue Telemetry Transport):一种基于发布/订阅(publish/subscribe)模式的轻量级通讯协议,通过订阅相应的主题来获取消息,是物联网(Internet of Thing)中的一个标准传输协议。该协议将消息的发布者(publisher)与订阅者(subscriber)进行分离,因此可以在不可靠的网络环境中,为远程连接的设备提供可靠的消息服务,使用方式与传统的MQ有点类似。TCP协议位于传输层,MQTT 协议位于应用层,MQTT 协议构建于TCP/IP协议上,也就是说只要支持TCP/IP协议栈的地方,都可以使用MQTT协议。为什么要用 MQTT协议?MQTT协议为什么在物联网(IOT)中如此受偏爱?而不是其它协议,比如我们更为熟悉的 HTTP协议呢?首先HTTP协议它是一种同步协议,客户端请求后需要等待服务器的响应。而在物联网(IOT)环境中,设备会很受制于环境的影响,比如带宽低、网络延迟高、网络通信不稳定等,显然异步消息协议更为适合IOT应用程序。HTTP是单向的,如果要获取消息客户端必须发起连接,而在物联网(IOT)应用程序中,设备或传感器往往都是客户端,这意味着它们无法被动地接收来自网络的命令。通常需要将一条命令或者消息,发送到网络上的所有设备上。HTTP要实现这样的功能不但很困难,而且成本极高。具体的MQTT协议介绍和实践,这里我就不再赘述了,大家可以参考我之前的两篇文章,里边写的也都很详细了。MQTT协议的介绍我也没想到 springboot + rabbitmq 做智能家居,会这么简单MQTT实现消息推送未读消息(小红点),前端 与 RabbitMQ 实时消息推送实践,贼简单~Websocketwebsocket应该是大家都比较熟悉的一种实现消息推送的方式,上边我们在讲SSE的时候也和websocket进行过比较。WebSocket是一种在TCP连接上进行全双工通信的协议,建立客户端和服务器之间的通信渠道。浏览器和服务器仅需一次握手,两者之间就直接可以创建持久性的连接,并进行双向数据传输。springboot整合websocket,先引入websocket相关的工具包,和SSE相比额外的开发成本<!-- 引入websocket --><dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId></dependency>服务端使用@ServerEndpoint注解标注当前类为一个websocket服务器,客户端可以通过ws://localhost:7777/webSocket/10086来连接到WebSocket服务器端@Component@Slf4j@ServerEndpoint("/websocket/{userId}")public class WebSocketServer { //与某个客户端的连接会话,需要通过它来给客户端发送数据 private Session session; private static final CopyOnWriteArraySet<WebSocketServer> webSockets = new CopyOnWriteArraySet<>(); // 用来存在线连接数 private static final Map<String, Session> sessionPool = new HashMap<String, Session>(); /** * 公众号:程序员小富 * 链接成功调用的方法 */ @OnOpen public void onOpen(Session session, @PathParam(value = "userId") String userId) { try { this.session = session; webSockets.add(this); sessionPool.put(userId, session); log.info("websocket消息: 有新的连接,总数为:" + webSockets.size()); } catch (Exception e) { } } /** * 公众号:程序员小富 * 收到客户端消息后调用的方法 */ @OnMessage public void onMessage(String message) { log.info("websocket消息: 收到客户端消息:" + message); } /** * 公众号:程序员小富 * 此为单点消息 */ public void sendOneMessage(String userId, String message) { Session session = sessionPool.get(userId); if (session != null && session.isOpen()) { try { log.info("websocket消: 单点消息:" + message); session.getAsyncRemote().sendText(message); } catch (Exception e) { e.printStackTrace(); } } }}前端初始化打开WebSocket连接,并监听连接状态,接收服务端数据或向服务端发送数据。<script> var ws = new WebSocket('ws://localhost:7777/webSocket/10086'); // 获取连接状态 console.log('ws连接状态:' + ws.readyState); //监听是否连接成功 ws.onopen = function () { console.log('ws连接状态:' + ws.readyState); //连接成功则发送一个数据 ws.send('test1'); } // 接听服务器发回的信息并处理展示 ws.onmessage = function (data) { console.log('接收到来自服务器的消息:'); console.log(data); //完成通信后关闭WebSocket连接 ws.close(); } // 监听连接关闭事件 ws.onclose = function () { // 监听整个过程中websocket的状态 console.log('ws连接状态:' + ws.readyState); } // 监听并处理error事件 ws.onerror = function (error) { console.log(error); } function sendMessage() { var content = $("#message").val(); $.ajax({ url: '/socket/publish?userId=10086&message=' + content, type: 'GET', data: { "id": "7777", "content": content }, success: function (data) { console.log(data) } }) }</script>页面初始化建立websocket连接,之后就可以进行双向通信了,效果还不错转载自https://www.studyjava.cn/post/2004
-
各位专家是否有使用java做mqtt服务器的指导案例,想让小熊派的数据上报上来并且展示
-
用华为云-设备接入的MQTT协议,不知道能否做设备间的语音对讲功能?若可以,请讲讲思路吧。
-
我在使用上层应用控制下层ESP32,需要用到平台下发命令,我是用Postman调测通过,但是使用Python时却报了{"error_code":"IOTDA.000001","error_msg":"Internal server error."}以下是我的代码import requests url = "XXXX" payload = { "service_id" : "SmokeDetectorControl", "command_name" : "ON_OFF", "paras" : { "value" : "1" } } headers = { 'x-auth-token': 'XXXXXX' } response = requests.request("POST", url, headers=headers, data=payload) print(response.text)与API explorer示例代码基本一致,希望哪位朋友能帮我看一下,不胜感激
-
作者 | Jim Meyers通过升级新的监控和数据采集(SCADA)系统,可再生能源公司在运营和业务管理方面获得更多优化和改进。使系统更易于满足合规要求的同时,还可以获得其它好处,对制造商来说,这是一个真正的双赢。只要问问Roeslein 可再生能源公司(RAE)就明白了,该公司投入了一套新的监控和数据采集系统(SCADA),除了合规性, 还在很多其它领域获得了改进。总部位于美国圣路易斯的可再生能源公司RAE, 拥有运营和开发沼气洗涤设施, 可将农业和工业生物废物转化为可再生天然气和可持续的副产品。系统集成商Roeslein&Associates(同属RAE 集团)为其升级了新的SCADA 系统,并使用了几种现代技术,包括云服务平台、SCADA 平台、边缘计算和消息队列遥测传输(MQTT)功能。 新系统可以从不同位置可靠地收集数据,将其存储在同一个地方,提供给用户访问。该系统帮助RAE 公司满足政府的合规性要求,也可以帮助其它公司满足这些要求。安全网络可以支持增加更多的公司,每个公司只能访问它们自己的数据。该系统提供了一种简化的方式来可视化数据、评估运营状态,并生成必要的月度报告,以确保合规性。 该系统提供了单一的事实来源,并简化了数据提交过程。RAE 公司每天还使用这个新平台来管理和优化生产,并能持续添加新的功能。检查所有阀门,确认是否有某个阀门开裂或泄漏。集成SCADA、云和边缘技术“连接所有部件并不困难。”负责该项目实施的系统工程师Mitchell Leefers 表示, 使用Amazon 云服务资源通常需要一段学习曲线,但通过边缘技术和MQTT,一切都可以完美地连接在一起。该项目涉及两个独立的数据采集系统。其中一个可以处理合规报告所需的数据,并可以以安全的方式将其它公司包括在内。该系统拥有一个虚拟的私有云。它还包括与可编程逻辑控制器(PLC)直接连接的以太网,以确保可靠的边缘设备通信。数据通过MQTT 传输。第二个系统仅用于RAE 公司收集其所有过程数据。这些数据被用于整体设备效率(OEE)、故障排除、维护跟踪和过程改进。 “ 当决定将其分为两个系统时,我们并没有现成的、在以前项目中使用过的任何模型,” Leefers 说,“我们做这个决定,完全是根据客户提供的信息和我们过去的经验。” SCADA 系统帮助RAE 公司向政府许可机构报告合规性。RAE 的财务经理Ivailo Chervenkov 说:“它通过用户友好的界面提供最新的运营数据,极大地促进了我们的合规工作。它还使我们可以以更高效的方式,来简化数据收集和分发。在使用该系统之前,所有数据都是通过繁琐的手动方法收集、过滤和分发的。” 现在,RAE 公司可以利用这些数据来实现更多功能, 以进一步改善运营。Cher venkov 说:“该系统使我们能够扩大报告工作,以应对不断增加的容量,以及更多的设施和增加的复杂性所带来的挑战。数据也变得更可靠和可用,加快了关键生产和预测活动。该系统帮助我们轻松确定出问题的领域,并快速找到实时解决方案,从而提高信息和决策的效率和质量。” SCADA 软件提供的不仅仅是合规性RAE 公司工程副总裁Eric Bancks 说:“SCADA 软件与我们的每日报告交互,允许单点访问数据。它提供了访问和可视化历史过程数据的工具,用于优化和故障排除的过程分析。我们能够绘制历史沼气产量与沼气预测模型对比图,并用经验数据修改预测模型,跟踪甲烷回收和甲烷不平衡。数据则用于识别‘可疑’仪表,以加快维修。” 新的SCADA 系统在解决问题的同时还节省了大量时间。Bancks 说:“因为大多数网点都离我们办公室很远。现在,我们可以在办公室收集数据并进行分析,为现场提供指导,而无需单程高达四五个小时的差旅。” 利用SCADA 平台中的分析工具,对数据进行趋势分析还会带来其它好处。Bancks 说:“我们的过程团队广泛使用定制趋势工具。它使我们的过程团队可以在同一个图上绘制多个变量,还可以将数据下载为CSV 文件,以供进一步操作和分析。”自定义趋势工具,具有不同的数据查询,如分钟快照、分钟平均值、小时平均值等。这有助于以任务所需的任何频率管理数据。RAE 已经预见到更好、更多地访问数据所带来的巨大价值,这对决策过程和预算编制至关重要。它提供了一致、准确和及时的数据,可供整个公司的不同员工日常使用。它还允许不同级别的管理层使用不同类型的信息,从而有助于改善决策。合规记录可帮助不同需求的企业创建合规记录的软件,也可以帮助没有合规要求的组织。“这类系统可用于任何类型的数据采集和分析,” Leefers 说,“我们已经建立了一个高可用性和高冗余的系统,通过互联网与中央数据库连接,可靠地从世界任何地方的运营中获取数据。这是我们解决的问题,但它也可能是其它不同场景的解决方案。” 系统集成商正在为该系统增加新的功能。“我们正在整合条形码系统, 以便更高效地跟踪天然气运输。” Leefers 说。许多生产场所没有与天然气管道的直接连接,所以在给卡车加油时,需要将天然气送到卸载站注入管道。每批货物都必须用一份文件进行跟踪,其中包含装载在卡车上的天然气信息。虽然,可以手动完成该过程的大部分工作,但需要大量运行人员手动输入数据以创建文档,且人工输入容易出错。新的条形码系统将使追踪每批气体及其数据变得更加容易和可靠。本文来自于控制工程中文版(CONTROL ENGINEERING China)2022年5月刊《应用案例》栏目:升级SCADA 系统:提升合规性和运营效率
-
【功能模块】小熊派的nb卡损坏了,开源社区里的实验手册都是利用nb-iot上云,学习了mqtt的智能农业实例后,想利用mqtt协议做智能烟感的实例【操作步骤&问题现象】个人能力有限,想知道如何修改mqtt示例中的哪些文件可以实现移植烟感模块用mqtt上云的功能【截图信息】【日志信息】(可选,上传日志内容或者附件)
yd_229306989
发表于2022-07-23 18:17:29
2022-07-23 18:17:29
最后回复
yd_229306989
2022-07-25 18:05:20
453 11 -
直接下载的Python Demo 文件,打开就报错,求指导。
-
【功能模块】D11_iot_cloud_oc_infrared目录下iot_cloud_oc_sample.c【操作步骤&问题现象】1、如果我不想连接到华为lot平台,而是其他mqtt服务器平台,这些参数应该填什么以及另外需要修改什么【截图信息】【日志信息】(可选,上传日志内容或者附件)
-
AR502H可以通过mqtt命令直接控制/dev/do0吗? 容器里已经有esdk和mosquitto运行了,向哪个topic地址publish命令可以控制/dev/do0? 软件二次开发里没有写这个内容,这个对开发APP很重要
-
【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
-
【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
上滑加载中
推荐直播
-
华为云码道Agent集成与鸿蒙实战2026/08/11 周二 19:00-21:00
王一男-华为云码道产品规划专家;李炎-华为云码道产品专家;彭江敏-华为云鸿蒙端云一体化开发专家
本次直播带你解读华为云码道7月份产品新特性、新功能。更有专家演示码道Agent Space × 钉钉机器集成实战,从0到1打通消息通道;码道鸿蒙端云一体化实战,快速搭建员工签到系统。
回顾中 -
华为云开发者AI素养直播课·第五期2026/09/04 周五 16:00-18:00
林华鼎-华为云AI开发者运营负责人;蒋春阳-华为云AI开发者案例开发专家
本期直播内容: AI工具体验营 · 第5-8课连讲。Agent-Team 多智能体协作完成毕业设计实践
回顾中 -
华为云开发者AI素养ClassRoom·第六期2026/09/08 周二 19:00-20:00
樊渊-2026华为软件挑战赛冠军
高手来了:看软挑高手解析二维排样问题—从工业难题到算法突破
回顾中
热门标签