ARTICLE DETAIL

资讯详情

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

SpringCloud微服务MQTT架构:设备消息统一接入、业务分发设计

SpringCloud微服务MQTT架构:设备消息统一接入、业务分发设计 SpringCloud微服务MQTT架构设备消息统一接入、业务分发设计作者黒漂技术佬上一篇文章我们搞定了 SpringBoot 单机版整合 MQTT。但如果你的智慧农业平台接入了几千个大棚、上万台设备天天吐着海量传感器数据——单机应用迟早会被撑爆。这个时候就需要微服务架构来拆解压力。本文带你设计一套基于 SpringCloud 的 MQTT 设备消息统一接入和业务分发方案。一、为什么要用微服务先说一个残酷的现实MQTT 消息处理和业务处理本质上是两种不同性质的负载。接入层IO 密集型大量 TCP 连接消息转发。瓶颈在网络和连接数。业务层CPU 密集型/IO 密集型数据解析、计算、入库。瓶颈在数据库和计算资源。把它们硬塞在一个进程里单体架构就会互相拖累。接入层被海量消息打满线程池的时候业务处理也跟着卡死。反之一个复杂的聚合查询把数据库拖慢可能影响到消息的正常接收。微服务的核心价值就是「各管各的独立伸缩」┌──────────────┐ │ MQTT Broker │ │ (EMQX) │ └──────┬───────┘ │ MQTT协议 ▼ ┌────────────────────────┐ │ device-gateway │ │ (接入网关 - 可横向扩展) │ └───────────┬────────────┘ │ RocketMQ / Kafka ▼ ┌────────────────┼────────────────┐ ▼ ▼ ▼ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │data-process │ │alert-service│ │control-svc │ │数据处理 │ │告警服务 │ │控制服务 │ └─────────────┘ └─────────────┘ └─────────────┘二、微服务职责划分来看看每个服务分别该干什么2.1 接入网关服务device-gateway这是整个系统的咽喉。所有设备消息先到这里再转发到内部系统。核心职责三板斧接收 MQTT 消息订阅所有设备数据 Topic消息转换MQTT 报文 → 统一内部消息体消息路由根据 Topic 下发到不同的 RocketMQ TopicSlf4jServicepublicclassMessageRoutingService{AutowiredprivateRocketMQTemplaterocketMQTemplate;// Topic → RocketMQ Topic 映射关系privatestaticfinalMapPattern,StringROUTE_MAPMap.of(Pattern.compile(agriculture/./sensor/.*),sensor-data,Pattern.compile(agriculture/./status),device-status,Pattern.compile(agriculture/./alarm),device-alarm);publicvoidroute(StringmqttTopic,Stringpayload){for(Map.EntryPattern,Stringentry:ROUTE_MAP.entrySet()){if(entry.getKey().matcher(mqttTopic).matches()){StringrocketTopicentry.getValue();// 构建统一消息体InternalMessagemsgInternalMessage.builder().mqttTopic(mqttTopic).payload(payload).timestamp(System.currentTimeMillis()).gatewayId(gatewayId)// 标识来自哪个网关实例.build();rocketMQTemplate.convertAndSend(rocketTopic,msg);return;}}log.warn(未匹配路由规则消息丢弃: {},mqttTopic);}}你可能注意到这里有个gatewayId——它用于标识消息来自哪个网关实例方便排查问题和负载追踪。2.2 设备管理服务device-service这个服务不处理传感器数据它管的是设备的「身份证」设备注册、激活设备首次入网时的握手流程设备认证连接 MQTT 时的用户名密码校验设备状态管理在线/离线/休眠设备OTA升级EMQX 这样的企业级 Broker 支持 HTTP 回调做设备认证——设备连接时 EMQX 会调你的 HTTP 接口你返回「允许」或「拒绝」即可。2.3 数据处理服务data-process从 RocketMQ 消费sensor-data消息解析后写入时序数据库Slf4jServiceRocketMQMessageListener(topicsensor-data,consumerGroupdata-process-group,selectorExpression*)publicclassSensorDataConsumerimplementsRocketMQListenerInternalMessage{AutowiredprivateInfluxDBServiceinfluxDBService;OverridepublicvoidonMessage(InternalMessagemsg){SensorDatadataparseSensorData(msg.getPayload());if(!validate(data)){log.warn(数据校验失败丢弃: {},data);return;}// 写入InfluxDB时序数据库influxDBService.write(data);}/** * 数据校验温度 -40℃ ~ 80℃湿度 0 ~ 100% */privatebooleanvalidate(SensorDatadata){if(data.getTemperature()-40||data.getTemperature()80){returnfalse;}if(data.getHumidity()0||data.getHumidity()100){returnfalse;}returntrue;}}2.4 告警服务alert-service与控制服务control-service告警服务从 RocketMQ 消费数据和预设阈值对比超出就发告警通知短信、钉钉、微信等。控制服务负责下发指令到设备。它不走消息队列而是通过 HTTP 直接调网关或者通过独立 MQTT 出站通道发送这样保证指令下发的低延迟。三、双层消息架构MQTT Broker 内部MQ这里要重点解释一下为什么我们需要两层消息队列MQTT Broker (设备 ↔ 服务器) │ ▼ RocketMQ / Kafka (服务 ↔ 服务)第一层 MQTT Broker是设备和服务器之间的桥梁。它解决了物联网最底层的问题轻量、省电、支持弱网、海量连接。第二层 RocketMQ / Kafka是服务与服务之间的桥梁。它解决了微服务架构中的问题解耦、削峰、异步、重试、死信、顺序消费。MQTT 是用来「接设备的」内部消息队列是用来「拆业务的」**。**两者各司其职。那能不能直接用 MQTT Broker 做服务间通信技术上可以但不推荐。MQTT 是按 Topic 发布订阅的无法提供 RocketMQ/Kafka 那种强大的消费组、消息回溯、Tag 过滤、事务消息等能力。四、Nacos 服务注册与发现既然是 SpringCloud 微服务当然少不了服务注册中心。这里用 Nacos# device-gateway 的配置spring:cloud:nacos:discovery:server-addr:127.0.0.1:8848namespace:smart-agriculturegroup:DEFAULT_GROUPapplication:name:device-gateway网关需要知道数据处理服务的地址吗不需要——它们通过 RocketMQ 异步通信完全解耦。但控制服务需要通过 Feign 调网关发指令时就需要 NacosFeignClient(namedevice-gateway)publicinterfaceDeviceGatewayClient{PostMapping(/api/command/send)ResultBooleansendCommand(RequestBodyCommandRequestrequest);}五、流量控制别让设备把服务打垮几千台设备同时发消息接入网关的压力可想而知。万一某批设备程序 Bug 死循环发消息直接把网关打挂了怎么办用 Sentinel 做流量控制Slf4jComponentpublicclassMqttMessageHandler{ServiceActivator(inputChannelmqttInputChannel)SentinelResource(valuemqtt-message-handle,blockHandlerhandleBlock)publicvoidhandleMessage(Message?message){Stringtopic(String)message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC);// 正常处理逻辑...}/** * 流量限制回调消息太多时触发 */publicvoidhandleBlock(Message?message,BlockExceptionex){log.warn(MQTT 消息处理被限流Topic: {},message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC));// 可以选择记录到本地缓冲区稍后处理// 或者直接丢弃传感器数据具有一定的时效性过期数据价值有限}}Sentinel 可以按 QPS每秒请求数或并发线程数限流。对于传感器数据建议按 QPS 限制比如单实例处理上限 5000 QPS超出就触发限流。六、智慧农业微服务拆分实战总结最终的服务划分和数据库对应关系微服务所属层数据库核心职责device-gateway接入层无状态MQTT桥接、消息路由、限流device-service业务层MySQL设备注册、认证、状态管理data-process数据层InfluxDB传感器数据解析、清洗、入库alert-service业务层MySQL阈值检测、告警通知control-service业务层MySQL指令下发、联动控制statistics-service数据层MySQL InfluxDB数据聚合、报表生成数据库的划分遵循「谁拥有数据谁掌管数据库」原则。device-service独享设备表其他服务要查设备信息必须通过 API 调它——这叫数据所有权是微服务设计中最容易被忽视但最重要的原则之一。总结微服务 MQTT 的本质是「接入和业务分离」。MQTT Broker 扛连接的RocketMQ/Kafka 扛业务的Nacos 管发现的Sentinel 管流控的。各自归位各司其职。这种架构可以轻松支撑 10 万 设备同时在线——接入网关加实例就行数据处理服务加消费者就行互不影响。
返回列表