ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

若依框架整合EMQX:MQTT消息收发与设备指令下发实战

若依框架整合EMQX:MQTT消息收发与设备指令下发实战 有朋友在群里问若依框架下怎么接EMQX设备上报的消息怎么收平台命令怎么下发到设备。这个问题我刚好从头到尾折腾过一遍从EMQX部署到Java侧订阅发布再到和若依的权限体系、业务模块打通踩了不少坑也沉淀了一套可以直接抄作业的方案。这篇就把整个思路和核心代码完整拆出来适合正在用若依做物联网后台、或者想把MQTT能力塞进Spring Boot项目的同学参考。先说清楚这套组合要解决什么问题。EMQX是业界很能打的开源MQTT消息服务器设备端或者客户端通过MQTT协议跟它通信Java这边需要一个MQTT客户端库负责订阅主题、接收消息、发消息若依则是国内用的非常多的Java后台快速开发框架基于Spring Boot自带了权限、用户、菜单、代码生成这些基础能力。把三者接起来场景就很典型了智能硬件设备通过MQTT上报数据若依后台负责接收解析、存储展示、下发控制指令或者前台系统通过订阅实时接收业务消息再通过若依的接口把指令发出去。1. 先想清楚为什么要用EMQX搞消息收发1.1 这个组合解决什么问题很多人在若依项目里遇到“设备上报数据”或者“服务端主动推送消息”的需求第一反应是写WebSocket。但WebSocket适合浏览器到服务器的实时通信如果设备端是单片机、ESP32、Android小程序这类终端WebSocket的兼容性和轻量程度都不如MQTT。MQTT是基于发布/订阅模式的轻量级消息协议特别适合网络不稳定、设备资源有限的物联网场景。EMQX作为MQTT broker负责消息的转发和主题管理。设备只需要连接到EMQX往某个主题发布消息Java服务端订阅这些主题就能实时收到数据Java服务端往设备专属主题发消息设备就能收到指令。所以用了这个组合之后各端角色就很清晰了设备端/终端通过MQTT客户端连接EMQX发布状态数据、接收控制指令EMQX消息中转站处理连接、主题匹配、QoS、离线消息Java服务端既是订阅者接收设备上报也是发布者下发指令再结合若依的鉴权、业务逻辑处理数据1.2 技术选型EMQX、Java客户端、若依各自扮演什么角色EMQX的部署方式很多最简单的是Docker一条命令就能跑起来控制台能实时看到连接数和消息流向。生产环境如果用集群EMQX也支持分布式扩展但那是后话先单机跑通再说。Java客户端方面大部分人选Eclipse Paho这应该是Java生态里最成熟的MQTT客户端库Spring Boot项目里直接加依赖就能用。有人问为什么不用EMQX官方Java客户端因为EMQX现在主推的emqx-extension-hook以及各类SDK本质上很多也是依赖Paho做底层的对于“订阅发送”这种常规需求Paho足够稳而且社区案例多、出了问题好查。若依这边无论你用的是RuoYi-Vue前后端分离版还是RuoYi单体版核心都是一个Spring Boot工程。我们需要做的就是把MQTT客户端的初始化和业务逻辑融入Spring的容器管理让MqttClient成为全局单例再配合若依的Controller、Service层做业务闭环。注意这套方案核心是“接到若依里”所以不是单纯写一个Java类连MQTT就完了而是要处理好Spring生命周期、配置管理、并发线程、消息分发这些工程化问题。2. 环境准备先把EMQX用Docker跑起来2.1 Docker安装EMQX与端口说明我建议直接用Docker部署EMQX干净、好卸载、版本切换方便。如果你用的是5.x版本一条命令就够了docker run -d --name emqx \ -p 1883:1883 \ -p 8083:8083 \ -p 8084:8084 \ -p 8883:8883 \ -p 18083:18083 \ emqx/emqx:5.6.0端口这块很容易搞混我整理了一张表照着用就行端口协议/用途1883MQTT TCP 端口Java客户端和设备端连这个8883MQTT SSL/TLS 端口生产环境加密连接用8083WebSocket端口浏览器端MQTT连接用8084WebSocket over TLS 端口18083EMQX Dashboard 控制台端口Web管理界面启动之后浏览器访问http://服务器IP:18083默认账号admin默认密码public登录后就能看到连接统计、订阅关系、消息速率这些信息。我第一次装完就是靠这个控制台确认broker是否有问题非常直观。如果你没有服务器只想在本机调试直接把服务器IP换成localhost就行。有人喜欢用docker-compose管理也完全OK核心配置就是上面这些端口映射。2.2 控制台验证与安全配置装完别急着写代码先用MQTTX这个桌面客户端工具连一下EMQX确认broker本身没问题。MQTTX是EMQX官方出的MQTT调试工具填上tcp://localhost:1883、客户端ID、用户名密码连接成功后就能手动发消息、订阅主题。这一步看起来多余其实特别重要——它能帮你在写Java代码之前排除掉“broker部署有问题”这种低级错误。还有一点必须提醒默认的admin/public账号在生产环境一定要改掉。EMQX控制台里可以配置新的用户名密码或者干脆做禁用默认账号。如果你是Java服务端作为唯一连接方也可以给EMQX创建一个专用账号只允许这个账号连接设备端再用独立的账号体系方便后续做权限隔离和审计。3. Java客户端接入订阅与发送的核心代码3.1 客户端库选型与依赖引入Paho的Maven坐标如下直接加到pom.xml里dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.2.5/version /dependency这里多说一句为什么不用MqttAsyncClient。Paho提供了同步MqttClient和异步MqttAsyncClient两个类。日常项目里很多人觉得异步性能好一上来就选MqttAsyncClient结果回调逻辑写得很别扭。实际上MqttClient内部也通过消息分发线程处理消息对于物联网后台这种场景单机每秒处理几百上千条消息完全扛得住而且代码写法更直观。如果想追求更高吞吐再考虑异步客户端不迟。3.2 连接参数与鉴权配置创建客户端之前先把连接参数想清楚。核心参数就这么几个brokerEMQX地址例如tcp://localhost:1883生产环境建议走SSLclientId客户端唯一标识同一时刻同一个clientId只能有一个连接重复会互相踢掉userName/passwordEMQX里配置的账号密码cleanSession是否清除会话要收离线消息就设falsekeepAliveInterval心跳间隔建议60秒connectionTimeout连接超时建议30秒automaticReconnect断线自动重连必须开启示例代码如下import org.eclipse.paho.client.mqttv3.MqttClient; import org.eclipse.paho.client.mqttv3.MqttConnectOptions; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; String broker tcp://localhost:1883; String clientId ruoyi-server- System.currentTimeMillis(); MemoryPersistence persistence new MemoryPersistence(); MqttClient client new MqttClient(broker, clientId, persistence); MqttConnectOptions options new MqttConnectOptions(); options.setUserName(emqx_admin); options.setPassword(your_password.toCharArray()); options.setCleanSession(false); options.setConnectionTimeout(30); options.setKeepAliveInterval(60); options.setAutomaticReconnect(true); // 遗嘱消息客户端异常掉线时broker代为发布一条消息 options.setWill(device/ruoyi/status, offline.getBytes(StandardCharsets.UTF_8), 1, true); client.connect(options);clientId这里我加了时间戳主要是防止多个实例启动时clientId冲突。如果你有多台服务器部署同一个服务一定要保证clientId唯一否则A连上了B被踢问题非常隐蔽。3.3 订阅主题与回调处理连接成功之后就是订阅。订阅之前先设计好主题规范。我常用的是这种分层结构device/{deviceId}/status // 设备状态上报 device/{deviceId}/telemetry // 设备遥测数据 command/{deviceId}/control // 平台下发控制指令主题里用{deviceId}做动态层订阅端可以用MQTT通配符一次性订阅一批主题。匹配单层#匹配多层所以订阅所有设备的状态用device//status订阅所有层级用device/#。订阅代码client.subscribe(device//status, 1); client.subscribe(device//telemetry, 1);回调接口是最核心的部分。Paho的MqttCallback里有三个方法分别对应断线、消息到达、消息发送完成client.setCallback(new MqttCallback() { Override public void connectionLost(Throwable cause) { // 自动重连开启后这里一般只做日志记录 log.error(MQTT连接断开, cause); } Override public void messageArrived(String topic, MqttMessage message) throws Exception { String payload new String(message.getPayload(), StandardCharsets.UTF_8); log.info(收到消息, topic{}, payload{}, topic, payload); // 这里先别写太重的业务逻辑后面会讲怎么分发 } Override public void deliveryComplete(IMqttDeliveryToken token) { // 消息发布完成回调一般用于确认 } });这里有个非常关键的细节messageArrived是在Paho内部的消息线程里执行的如果你在这个方法里直接操作数据库、调用远程接口很容易把线程阻塞住导致后续消息处理延迟。正确做法是收到消息后立即交给业务线程池异步处理后面第4节会讲。3.4 消息发送与QoS选择发送消息用publish方法String topic command/ deviceId /control; String payload {\type\:\reboot\,\timestamp\:1700000000}; MqttMessage message new MqttMessage(payload.getBytes(StandardCharsets.UTF_8)); message.setQos(1); message.setRetained(false); client.publish(topic, message);QoS的选择值得单独说。QoS级别语义场景0最多一次可能丢消息实时性高、允许丢失的遥测数据1至少一次可能重复控制指令、状态上报配合去重2恰好一次性能开销大计费、关键日志等极重要消息我实际中大多数场景用QoS 1就够。注意QoS 1存在重复消息的可能如果业务上不能容忍重复可以在Java层用消息ID做去重。另外setRetained(true)是保留消息新订阅的客户端能立刻收到最近一次消息适合设备状态这类场景但别滥用不然每次设备上线都会收到一堆历史状态。4. 与若依框架整合把MqttClient变成Spring Bean4.1 配置管理与初始化销毁如果只是写一个独立类问题不大但要融入若依这种有严格分层和配置规范的项目最好把它做成Spring管理的单例Bean同时把连接参数抽到application.yml里。先加配置mqtt: broker: tcp://localhost:1883 client-id: ruoyi-server username: emqx_admin password: your_password default-topics: - device//status - device//telemetry qos: 1 clean-session: false然后写一个配置类Configuration ConfigurationProperties(prefix mqtt) Data public class MqttProperties { private String broker; private String clientId; private String username; private String password; private ListString defaultTopics; private int qos; private boolean cleanSession; }再写一个MqttConfig负责创建client、连接、订阅Configuration RequiredArgsConstructor Slf4j public class MqttConfig { private final MqttProperties mqttProperties; Bean(destroyMethod close) public MqttClient mqttClient() throws MqttException { String clientId mqttProperties.getClientId() - System.currentTimeMillis(); MqttClient client new MqttClient(mqttProperties.getBroker(), clientId, new MemoryPersistence()); MqttConnectOptions options new MqttConnectOptions(); options.setUserName(mqttProperties.getUsername()); options.setPassword(mqttProperties.getPassword().toCharArray()); options.setCleanSession(mqttProperties.isCleanSession()); options.setConnectionTimeout(30); options.setKeepAliveInterval(60); options.setAutomaticReconnect(true); try { client.connect(options); for (String topic : mqttProperties.getDefaultTopics()) { client.subscribe(topic, mqttProperties.getQos()); log.info(订阅主题: {}, topic); } log.info(EMQX连接成功, broker{}, mqttProperties.getBroker()); } catch (MqttException e) { log.error(EMQX连接失败, e); throw e; } return client; } }destroyMethod close是重点Spring容器关闭时会自动调用MqttClient.close()释放连接不会把后台线程留在JVM里。如果你用的是旧版Spring Boot也可以手动写PreDestroy方法确保先断开再关闭连接。4.2 消息分发与业务解耦订阅回调拿到裸消息之后如果直接在回调里写if (topic.startsWith(device/))这种判断代码很快会变成一坨。我建议用Spring的事件机制做解耦让消息处理模块可以独立扩展。先定义事件public class MqttMessageEvent extends ApplicationEvent { private final String topic; private final String payload; public MqttMessageEvent(Object source, String topic, String payload) { super(source); this.topic topic; this.payload payload; } // getter... }在messageArrived里发布事件Component RequiredArgsConstructor Slf4j public class MqttCallbackHandler implements MqttCallback { private final ApplicationEventPublisher eventPublisher; Override public void messageArrived(String topic, MqttMessage message) { String payload new String(message.getPayload(), StandardCharsets.UTF_8); eventPublisher.publishEvent(new MqttMessageEvent(this, topic, payload)); } Override public void connectionLost(Throwable cause) { log.error(MQTT连接丢失, cause); } Override public void deliveryComplete(IMqttDeliveryToken token) { // 可做消息确认统计 } }然后写业务监听器在监听器里做主题解析和业务处理Component RequiredArgsConstructor Slf4j public class DeviceMessageListener { private final DeviceDataService deviceDataService; EventListener Async(mqttTaskExecutor) public void onMqttMessage(MqttMessageEvent event) { String topic event.getTopic(); String payload event.getPayload(); // 处理遥测数据 if (topic.matches(^device/[^/]/telemetry$)) { String deviceId topic.split(/)[1]; deviceDataService.handleTelemetry(deviceId, payload); } else if (topic.matches(^device/[^/]/status$)) { String deviceId topic.split(/)[1]; deviceDataService.handleStatus(deviceId, payload); } } }注意监听器上加的Async(mqttTaskExecutor)。我踩过一个大坑设备量一多消息并发上来之后回调线程被业务逻辑里的数据库操作拖死EMQX控制台里看到订阅连接正常但消息处理延迟越来越大。后来把所有业务处理扔到独立线程池问题立刻缓解。线程池配置放在若依的AsyncConfig里或者单独写一个配置类Configuration public class MqttThreadPoolConfig { Bean(mqttTaskExecutor) public Executor mqttTaskExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); executor.setMaxPoolSize(16); executor.setQueueCapacity(1000); executor.setThreadNamePrefix(mqtt-task-); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }CallerRunsPolicy的意思很简单线程池满了之后任务回退到调用线程执行不会直接丢弃消息。对我们这种消息型业务来说丢消息比慢更可怕。4.3 对外提供HTTP接口设备要接入平台肯定得开放接口。若依里最顺手的做法就是写一个Controller把发布消息能力封装成REST接口。先写一个ServiceService RequiredArgsConstructor Slf4j public class MqttPublishService { private final MqttClient mqttClient; public void publish(String topic, String payload, int qos) { try { MqttMessage message new MqttMessage(payload.getBytes(StandardCharsets.UTF_8)); message.setQos(qos); message.setRetained(false); mqttClient.publish(topic, message); log.info(发送MQTT消息, topic{}, payload{}, topic, payload); } catch (MqttException e) { log.error(发送MQTT消息失败, e); throw new ServiceException(指令下发失败); } } }Controller层就正常写配合若依的权限注解做接口保护RestController RequestMapping(/iot/command) RequiredArgsConstructor public class IotCommandController extends BaseController { private final MqttPublishService mqttPublishService; PostMapping(/send) PreAuthorize(ss.hasPermi(iot:command:send)) public AjaxResult send(RequestBody CommandSendRequest request) { String topic command/ request.getDeviceId() /control; mqttPublishService.publish(topic, request.getPayload(), 1); return success(); } }PreAuthorize是若依自带的权限控制这样指令下发接口就能纳入若依的菜单权限管理体系只有被授权的人才能用。4.4 把设备消息推到前端页面Java后台收到设备消息之后业务上通常会有一个动作把最新状态实时展示到Web页面上。这时候就要用到WebSocket了。若依本身没有内置WebSocket模块需要自己加。我采用的方式是后端监听MQTT消息入库后通过Spring WebSocket主动推送给前端登录用户。Spring WebSocket的配置不复杂核心就三个步骤引入spring-boot-starter-websocket依赖配置一个WebSocketConfigurer注册WebSocket handlerhandler里维护一个ConcurrentHashMap保存会话向前端推送消息消息推送这部分其实就是把MQTT收到的事件通过WebSocket转发出去。当然要注意前端页面离开后及时关闭会话否则连接数会无限增长。5. 常见问题与避坑实录5.1 IDEA导入若依多模块工程报“error adding module to project: null”这个报错在搜索相关词里热度很高很多人导入若依工程时就卡在这一步。实际原因大多是IDEA在识别多模块Maven工程时子模块没有被正确加载。解决办法是打开IDEA右侧Maven面板点击刷新按钮重新导入项目如果还是不行在项目根目录的pom.xml上右键选择Add as Maven Project让IDEA重新加载父POM基本都能解决。另外确认一下本地JDK版本和Maven配置若依官方要求的JDK版本和IDEA默认SDK不一致时也会出现这种奇奇怪怪的导入问题。5.2 客户端连接不稳定设备消息收不到先看EMQX Dashboard的“连接”页面确认客户端在线状态。如果发现客户端反复掉线重连第一嫌疑是clientId冲突——多个客户端实例用了同一个clientId后连接的把先连接的踢下线形成死循环。解决方法是保证clientId全局唯一或者在clientId后面拼上UUID/时间戳。如果连接稳定但收不到消息检查订阅主题和发布主题是否匹配。device/和device/#虽然都合法但语义不同只匹配一层#匹配多层中间差一个层级都收不到。5.3 离线消息没了设备上线看不到最近状态默认情况下EMQX只给当前在线的客户端转发消息。如果一个设备离线期间服务端发布了一条指令设备上线后是收不到的。要解决这就得靠两个手段一是刚才讲过把cleanSession设为false让broker保存会话和设备端的订阅关系二是EMQX 5.x支持会话过期时间在连接时可以通过Session Expiry Interval设置过期时长。不过这里提醒一句长期保留离线消息会占用broker内存生产环境建议评估设备数量后设置合理的过期时间比如一天或一周。5.4 消息乱码问题MQTT消息的payload本身就是字节数组没有强制编码。Java端发送时用payload.getBytes(StandardCharsets.UTF_8)接收时用new String(message.getPayload(), StandardCharsets.UTF_8)两端统一UTF-8就不会乱码。如果设备端用了GBK编码那Java这边就得配合设备协议做转码这个属于协议对接问题不在框架层面。5.5 高频消息把数据库打爆有些设备几秒上报一次数据几十台设备叠加起来数据库压力非常大。我的做法是收到消息后先写Redis缓存达到一定数量或者固定时间窗口后再批量落库。这个类似“攒批”的设计能大幅降低数据库的写压力。你可以在若依里直接用RedisTemplate配合定时任务实现不需要额外引入组件。5.6 EMQX容器重启后配置丢失Docker部署EMQX如果不挂载数据卷容器删除后所有配置、账号、数据都会消失。解决办法是在启动命令里加-v参数把EMQX的数据目录和配置目录挂载到宿主机。具体哪些目录需要挂载以官方镜像文档为准基础的做法是挂载/opt/emqx/data和/opt/emqx/etc两个目录这样容器销毁重建后数据和配置都还在。6. 一点实操心得这套东西我前后调了两天最深的感受是连接EMQX本身不难难的是把它塞进若依的工程体系之后各种工程化细节。比如MqttClient的生命周期管理、消息处理线程的动作、主题规范的设计这几个点如果一开始没想清楚后面越写越乱。如果你只是想把消息接到若依里先跑通建议按这个顺序动手先用Docker起EMQX → 用MQTTX手动收发验证 → 写一个最简单的Java客户端连接订阅 → 再把它改造成Spring Bean → 最后接业务逻辑。别一上来就写一整套代码出了问题你会分不清是broker的问题、客户端的问题还是业务的问题。另外主题命名规范一定要在一开始就定好最好写进团队文档里。我用的是device/{deviceId}/suffix这种结构设备属性、上下线、遥测、命令分得清清楚楚后面做权限控制、消息路由都方便。如果主题命名混乱后面扩展一个模块就得动一堆代码那才叫真正的灾难。
返回列表