SpringBoot集成MQTT:智能售货柜“取货即走”订单链路实战
发布时间:2026/10/1 5:11:34 作者:尧图编辑部 阅读量:1,286

去年在社区做智能售货柜试点的时候我踩了不少坑才把整套链路跑通。这个项目名字听起来挺玄乎——取货即走说白了就是用户扫码开门、拿走商品、关门自动扣款全程没有扫码支付这一步。后台的核心逻辑全靠SpringBoot搭的服务端设备和服务器之间的通信则完全依赖MQTT协议。当时团队里有人提议用HTTP轮询被我一票否决了后面的事实也证明这个决定是对的。这篇文章我会把整个项目的架构思路、SpringBoot集成MQTT的具体写法、取货即走的订单链路实现以及我在调试过程中遇到的典型问题全部梳理一遍。如果你正准备做无人零售、智能货柜、共享设备这类物联网场景或者是想学SpringBoot整合MQTT但找不到完整案例的开发者这篇文章应该能帮你少走不少弯路。1. 项目拆解取货即走场景下的技术选型1.1 取货即走的业务闭环先把这个项目的业务逻辑说清楚。智能售货柜不同于传统的无人售货机传统机器是你选一个格子投币或者扫码货掉下来整个过程是一步一步来的。取货即走模式更像无人便利店柜门上贴着二维码用户扫码后柜门解锁拿走想要的商品关上门的瞬间系统根据重量变化或者视觉识别结果自动算出拿了什么、扣多少钱。整个业务闭环拆开看是这样的用户扫码服务端下发开柜指令。柜门打开用户拿走商品。关门触发设备端传感器常见的是重力感应货架、RFID、或者摄像头视觉识别。设备端把取货前后的重量差异或识别结果上报给服务端。服务端根据差异匹配商品、生成订单、完成扣款。用户收到扣款通知交易结束。这里面最关键的环节是第4步和第5步。设备上报的数据格式、上报时机、上报可靠性直接决定了订单能不能算对。而SpringBoot服务端的职责就是接收设备消息、解析数据、匹配商品、状态流转、对接支付渠道。1.2 为什么是MQTT而不是HTTP轮询这是项目初期争得最凶的一个技术选型问题。有人说设备端每5秒调一次HTTP接口把状态上报上来不就完了干嘛要搞MQTT这种带连接状态的消息协议。我从三个维度分析一下这个问题的答案。第一是功耗和网络开销。售货柜通常分布在商场、地铁站、社区门口很多地方没有稳定的WiFi设备一般用4G模组。HTTP轮询意味着每次都要建立TCP连接、发送HTTP报文、等待响应这玩意儿对电池供电的设备来说是灾难。MQTT基于长连接一次连接建立后可以持续传输消息报文头最小只有2个字节省电省流量。第二是设备状态的实时可感知性。MQTT有遗嘱消息LWT机制设备异常掉线时Broker能立刻推送离线通知给服务端。HTTP轮询只能做到超时未上报就判定离线这个判定往往要等好几个轮询周期响应太慢。对于售货柜这种需要实时监控设备健康状态的场景MQTT的在线感知能力几乎是不可替代的。第三是发布订阅模型天然适合一对多的指令下发。服务端要向多台设备下发开柜指令如果走HTTP就是逐个调用设备接口如果走MQTT只需要向某个topic发布一条消息所有订阅了这个topic的设备都能收到。反过来多台设备上报数据服务端也只需要订阅一个统配的topic消息自动汇聚。想加一台设备不用改代码它上线后订阅自己的topic即可。MQTT协议层面QoS服务质量分0、1、2三级支付相关场景建议至少QoS 1保证消息不丢。这个点后面我会详细说。1.3 整体架构与模块划分整个系统的模块划分可以分为四块设备端基于ESP32或STM32主控配合重力传感器、电磁锁、4G/WiFi模组。设备端负责采集数据、订阅开柜指令、发布状态消息。MQTT Broker我用的是EMQX开源版够用。它负责消息的路由转发、设备认证、遗嘱消息处理。SpringBoot服务端整个项目的核心大脑。包含MQTT客户端、设备管理模块、订单模块、商品库存模块、支付模块。用户端微信小程序/App负责扫码、展示订单、支付结果回调。服务端内部我按照包结构划分了职责configMQTT配置、Redis配置、线程池配置mqtt消息监听器、消息模板、topic常量device设备注册、设备上下线管理order订单状态机、订单流转、超时任务product商品管理、库存扣减pay支付渠道对接、回调处理common统一返回、异常处理、工具类模块之间通过Spring的事件机制解耦。比如设备上线时发布DeviceOnlineEvent订单模块监听这个事件后把之前未完成的历史订单触发结算库存模块监听扣款成功事件扣减对应商品库存。用事件驱动的好处是新增一个业务模块比如后续加营销活动不需要改动已有模块的代码。2. SpringBoot服务端搭建从依赖到MQTT接入2.1 工程初始化与依赖选型项目用的是Spring Boot 2.7.xJDK 1.8。这个组合是目前生产环境最稳的搭配坑少资料多。当然如果你的项目是全新启动用Spring Boot 3.x也没问题只是下面的代码在javax到jakarta包名上会有差异。MQTT客户端的选型我对比过两个方案一个是Eclipse Paho的org.eclipse.paho.client.mqttv3另一个是Spring Integration的spring-integration-mqtt。Paho是纯客户端库灵活但需要自己管理重连、线程池Spring Integration把MQTT封装成了消息通道和Spring的ServiceActivator、MessagingGateway完美融合。我最终选了Spring Integration因为它的入站消息处理和出站发送可以像Spring MVC写接口一样简洁。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-integration/artifactId /dependency dependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3/version /dependency除了MQTT相关依赖Redis是必须的。设备在线状态、消息幂等、分布式锁、高频库存扣减这些都得靠Redis扛。MyBatis-Plus纯粹是为了开发效率没必要手写一堆单表CRUD。2.2 MQTT连接配置与客户端工厂项目的核心配置写在application.yml里。mqtt: broker: url: tcp://192.168.1.100:1883 username: vending password: vending-secret client: id: springboot-server-001 timeout: 30 keepalive: 60 topic: device-status: /device/{deviceId}/status device-event: /device/{deviceId}/event server-command: /server/{deviceId}/command这里重点说一下几个参数的含义。keepalive是心跳间隔单位秒。设备端和服务端都会按这个间隔发送PINGREQ包保持连接活跃。设太短会增加流量开销设太长会导致掉线检测迟钝。对于售货柜这类固定电源供电的设备60秒是一个稳妥的默认值。clientId必须全局唯一。MQTT协议要求每个客户端连接时使用独立的clientId如果两个客户端用了同一个IDBroker会踢掉前一个连接。这个也是生产环境最常见的连接异常原因之一后面我会单独说。连接工厂用DefaultMqttPahoClientFactory它是Spring Integration MQTT的核心工厂类负责创建真实的MQTT连接。Configuration public class MqttConfig { Value(${mqtt.broker.url}) private String brokerUrl; Value(${mqtt.broker.username}) private String username; Value(${mqtt.broker.password}) private String password; Value(${mqtt.client.id}) private String clientId; Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory new DefaultMqttPahoClientFactory(); MqttConnectOptions options new MqttConnectOptions(); options.setServerURIs(new String[]{brokerUrl}); options.setUserName(username); options.setPassword(password.toCharArray()); options.setCleanSession(true); options.setKeepAliveInterval(60); options.setAutomaticReconnect(true); options.setMaxReconnectDelay(30000); factory.setConnectionOptions(options); return factory; } }setAutomaticReconnect(true)之后Paho客户端在网络断开时会自动重连不需要自己写重连线程。很多初学者不知道这个API自己辛辛苦苦写一套重连逻辑其实官方已经提供了。2.3 消息订阅入站消息处理Spring Integration MQTT的入站适配器MqttPahoMessageDrivenChannelAdapter负责订阅topic、接收消息、把消息转成Spring的Message对象。Bean public MessageProducer mqttInbound() { MqttPahoMessageDrivenChannelAdapter adapter new MqttPahoMessageDrivenChannelAdapter(clientId -inbound, mqttClientFactory(), /device//status, /device//event); adapter.setCompletionTimeout(5000); adapter.setQos(1); adapter.setOutputChannel(mqttInputChannel()); return adapter; } Bean public MessageChannel mqttInputChannel() { return new DirectChannel(); }/device//status里的是MQTT的通配符匹配任意一级。这样所有设备的状态消息都会汇聚到这个适配器里。消息过来之后通过ServiceActivator注解监听消息通道这是整个服务端接收设备消息的唯一入口。Component public class MqttMessageReceiver { private static final String PAYLOAD_HEADER mqtt_receivedPayload; private static final String TOPIC_HEADER mqtt_receivedTopic; ServiceActivator(inputChannel mqttInputChannel) public void handleMessage(Message? message) throws Exception { String topic message.getHeaders().get(TOPIC_HEADER, String.class); String payload message.getHeaders().get(PAYLOAD_HEADER, String.class); String deviceId parseDeviceId(topic); MqttMessage msg JSON.parseObject(payload, MqttMessage.class); String type msg.getType(); switch (type) { case status: deviceStatusService.handleStatusReport(deviceId, msg); break; case event: deviceEventService.handleDeviceEvent(deviceId, msg); break; default: log.warn(unknown message type: {}, type); } } private String parseDeviceId(String topic) { // /device/{deviceId}/status - {deviceId} return topic.split(/)[2]; } }这里我把消息分发做成了策略模式type字段区分消息类型设备上报的重量变化事件、心跳事件、故障事件分别走不同的处理器。好处是后续每新增一种设备消息类型只需要新增一个case分支和对应的Service不会影响已有逻辑。有个细节mqtt_receivedPayload这个header才是真正的消息内容直接用message.getPayload()拿到的是适配器包装过的对象。第一次接入的人很容易在这里卡住我当时的排查过程花了一个多小时最后打印出所有headers才发现问题所在。2.4 消息发布指令下发指令下发是从服务端向设备发消息比如开柜指令、重启指令、亮灯指令。Spring Integration MQTT提供MqttPahoMessageHandler作为出站适配器。Bean ServiceActivator(inputChannel mqttOutboundChannel) public MessageHandler mqttOutbound() { MqttPahoMessageHandler handler new MqttPahoMessageHandler(clientId -outbound, mqttClientFactory()); handler.setAsync(true); handler.setDefaultQos(1); handler.setDefaultTopic(/server/command); return handler; } Bean public MessageChannel mqttOutboundChannel() { return new DirectChannel(); }发送消息的时候我封装了一个MqttGateway组件业务方不需要关心底层MQTT细节。Component public class MqttGateway { Autowired private MessageChannel mqttOutboundChannel; public void sendCommand(String deviceId, String payload) { MessageString message MessageBuilder.withPayload(payload) .setHeader(mqtt_topic, /server/ deviceId /command) .setHeader(mqtt_qos, 1) .build(); mqttOutboundChannel.send(message); } }mqtt_topic这个header是用来动态指定发送目标的。每个设备的指令topic不同直接在header里设置比硬编码到handler里灵活得多。3. 核心链路实现从开柜取货到自动扣款3.1 设备上下线感知与在线状态管理设备管理是整个系统稳定运行的基础。设备掉线没人知道用户扫了码却开不了门这是最糟糕的体验。我通过MQTT的遗嘱消息和定期心跳消息实现设备在线状态管理。设备端在连接时设置遗嘱消息Will Message/* * 伪代码展示设备端的遗嘱设置逻辑 */ MqttConnectOptions options new MqttConnectOptions(); options.setWill(/device/{deviceId}/status, {\type\:\offline\}.getBytes(), 1, false);遗嘱消息的意思是如果设备异常断开比如断电、断网Broker会立即代设备向这个topic发布一条消息。服务端订阅了/device//status就能很快感知设备离线。设备正常上线时会主动发一条online消息之后每隔30秒发一条包含电量、温度、信号强度的心跳消息。服务端的DeviceStatusService维护一张设备状态表字段包括字段类型说明device_idvarchar设备编号statustinyint0离线 1在线 2故障last_online_timedatetime最近上线时间last_heartbeat_timedatetime最近心跳时间firmware_versionvarchar固件版本判断设备是否离线的逻辑是收到offline遗嘱消息时立即置为离线如果超过3个心跳周期没有收到心跳消息后台定时任务主动巡检并将其置为离线。在线状态下用户扫码后服务端会先检查设备状态。如果设备离线直接提示用户设备不可用避免出现扫码后开不了门的情况。3.2 订单状态机设计与异常补偿取货即走模式下的核心对象是订单。订单的创建时机不是关门后而是开柜时。用户扫码开柜的那一瞬间服务端就创建一笔初始状态的订单后续所有操作都围绕这笔订单流转。我把订单状态设计成6个状态枚举状态值含义CREATED0已创建柜门已开PICKED1用户已取货等待结算SETTLING2结算中正在匹配商品PAID3已扣款交易完成SETTLE_FAILED4结算失败需人工处理CANCELLED5已取消用户未取货状态流转是这样的用户扫码开柜 - 创建订单状态为CREATED。关门后设备上报重量数据 - 服务端收到上报状态为PICKED。服务端匹配商品并计算金额 - 状态为SETTLING。调用支付接口扣款成功 - 状态为PAID。如果扣款失败或者商品匹配不出来 - 状态为SETTLE_FAILED触发告警运维人工介入。状态流转我用了一个轻量级的状态机实现用Map维护每个状态允许跳转的目标状态集合每次更新状态时校验合法性。这个方法简单直接不需要引入状态机框架排查问题的时候也容易跟踪。public class OrderStateMachine { private static final MapInteger, SetInteger TRANSITIONS new HashMap(); static { TRANSITIONS.put(0, new HashSet(Arrays.asList(1, 2, 5))); TRANSITIONS.put(1, new HashSet(Arrays.asList(2, 4))); TRANSITIONS.put(2, new HashSet(Arrays.asList(3, 4))); TRANSITIONS.put(3, new HashSet()); TRANSITIONS.put(4, new HashSet()); TRANSITIONS.put(5, new HashSet()); } public static boolean canTransition(Integer from, Integer to) { SetInteger allowed TRANSITIONS.get(from); return allowed ! null allowed.contains(to); } }订单一旦进入SETTLE_FAILED状态系统会发钉钉告警给运营人员同时在用户端展示订单处理中请稍后查看的提示。我见过太多无人零售项目在这个环节处理不当用户货拿了钱没扣对客服被骂口碑直接崩塌。所以异常订单必须第一时间告警而不是默默失败。3.3 扣款流程与库存扣减的一致性扣款流程设计上遵循一个原则先算账、后扣款、再减库存。顺序不能反。关门后设备上报重量变化服务端的结算流程如下根据设备ID查询该设备关联的商品列表每个货架上有哪些商品、每个商品的重量阈值。将上报的重量变化值逐一和商品重量比对允许±3克的误差。匹配出商品明细列表计算总金额。调用微信/支付宝免密代扣接口扣款。扣款成功后扣减对应商品库存。这里最容易出错的是第2步如果货架上有两种重量接近的商品重量传感器没法精确区分。我的方案是多传感器交叉验证每个货架有四个称重传感器分别上报四个方向的重量组合成一条特征向量再和商品库里的标准特征向量做匹配准确率能到95%以上。如果匹配置信度低于阈值就把这笔订单置为SETTLE_FAILED宁可人工介入不瞎扣钱。库存扣减用的是Redis加数据库双写。Redis里存每个设备的可用库存扣减时用Lua脚本保证原子性-- 扣减库存脚本 local key KEYS[1] local qty tonumber(ARGV[1]) local current tonumber(redis.call(get, key)) if current nil or current qty then return -1 end redis.call(decrby, key, qty) return current - qty扣减成功后再异步更新MySQL里的库存表。如果MySQL更新失败会有定时任务做对账以Redis的流水记录为基准修正数据库库存。用Redis的原因是高频扣减场景下直接操作MySQL的行锁会成为瓶颈一台热门点位的高峰期扣减QPS能到几十MySQL撑得住但没必要硬扛。3.4 消息可靠性QoS、幂等与乱序智能售货柜场景里消息可靠性是服务端设计的重中之重。我总结出三个核心问题丢消息、重复消息、乱序消息。先说丢消息。MQTT有三种QoS级别QoS 0最多一次消息可能丢失。QoS 1至少一次消息保证送达但可能重复。QoS 2恰好一次消息不丢不重但开销最大。设备上报的结算数据直接关系到扣款我全部用QoS 1。QoS 2在4G弱网环境下会有明显的性能回退容易造成消息阻塞实际项目中很少用。再说重复消息。QoS 1下消息可能重复到达服务端如果不做幂等处理用户拿了一瓶可乐可能被扣两次钱。我的做法是每台设备上报消息时带一个自增序号seq服务端收到消息后先检查Redispublic boolean isDuplicate(String deviceId, long seq) { String key mqtt:msg: deviceId : seq; Boolean success redisTemplate.opsForValue() .setIfAbsent(key, 1, Duration.ofMinutes(5)); return !Boolean.TRUE.equals(success); }setIfAbsent是Redis的SETNX操作同一个seq的消息只有第一次能被写入后续的重复消息直接丢弃。这个方案的过期时间设置为5分钟因为MQTT的重复消息通常出现在网络抖动后的短时间内过期后序号可以重新使用。最后说乱序问题。设备上报的seq就是解决乱序的关键。服务端在Redis里给每台设备维护一个最近一次处理的seq当新消息的seq比当前值大1认为是正常顺序如果大于1说明中间有消息丢了触发补拉流程服务端向设备发送一条PULL指令设备把丢失序号区间内的缓存数据重新上报。这套机制组合起来基本能保证在69%丢包率的弱网环境下订单仍然能正确生成实测结果还不错。4. 调试与避坑从日志到线上实战4.1 开发调试利器MQTTX与mosquitto项目开发时最让我头疼的不是SpringBoot代码本身而是调试设备端和服务端的消息通信。没有趁手的工具每次都要两头打日志效率极低。后来我固定下来三件套MQTTX客户端一个跨平台的MQTT调试工具支持Windows/Mac/Linux。我可以手动模拟一台设备向服务端发各种格式的消息极大方便了联调。mosquitto_sub/mosquitto_pub命令行工具适合在服务器上快速验证Broker连通性。比如服务端发布了一条开柜指令我用mosquitto_sub -t /server/#就能直观看到指令有没有出来。EMQX DashboardWeb管理界面可以实时查看所有设备的连接状态、订阅关系、收发消息速率。排查设备频繁掉线问题时Dashboard的日志列表是最直接的证据。我推荐团队成员在开发时都用MQTTX模拟设备配合调试把设备端和服务端的联调前置到开发阶段而不是等到样机出来再联调。4.2 频发的幽灵订单问题排查项目上线第一周运营反馈了一个诡异的问题有些柜子没人使用后台却出现了扣款订单。幽灵订单这个词就是我们当时起的。排查过程是这样的查看订单日志发现这些订单的创建时间集中在凌晨3点到5点。查设备运行日志发现这个时间段的柜机上报过重量变化事件。到现场拆机检查发现夜间温度变化剧烈金属货架热胀冷缩产生的微弱形变让重力传感器产生了3~5克的静漂。根因找到了凌晨温度下降到一定阈值传感器零点漂移导致设备误判为有人取货于是自动上报了取货事件。我们的解决方案分两层硬件层在称重传感器支架上加一层缓冲垫减少热胀冷缩对传感器受力的影响。软件层增加最小变化阈值判断重量变化小于15克的事件直接忽略不触发结算流程。另外增加夜间结算抑制窗口凌晨2点到5点期间不上报非必要的重量事件但保留心跳。这个case给我最大的教训是物联网系统的坑往往藏在硬件和环境的边界上纯靠写代码解决不了所有问题一定要有现场排查的能力。4.3 弱网环境下的设备消息补偿商场地下室、地铁站台这些位置的4G信号很差设备经常处于连得上但经常断的状态。断线时用户正好在取货这个消息就很容易丢。我的方案是设备端本地缓存加服务端补拉。设备端每次上报消息时除了实时发送还会在本地Flash中缓存最近100条消息。服务端根据seq的跳变发现丢消息后通过MQTT向设备下发补拉指令{command:PULL,fromSeq:101,toSeq:105}。设备收到后把缓存中的101到105序号的消息重新按序上报。服务端补拉逻辑public void handlePulledMessages(String deviceId, ListMqttMessage messages) { // 按seq排序 messages.sort(Comparator.comparingLong(MqttMessage::getSeq)); for (MqttMessage message : messages) { long seq message.getSeq(); if (isDuplicate(deviceId, seq)) { continue; } processMqttMessage(deviceId, message); } }这套机制上线后消息丢失导致的问题基本清零。核心经验是在MQTT这种低带宽长连接的协议里不能指望链路永远稳定一定要在业务层设计补偿机制。设备端Flash缓存的空间很小所以只缓存重量事件和开关门事件这两类关键数据心跳这类非关键消息丢了也就丢了。4.4 上线前必须做的几件事如果让我重新做一遍这个项目上线前我一定会更认真地做这几项检查第一clientId唯一性检查。SpringBoot服务端连了一个clientId如果有多套环境共用同一个Broker两个服务用同一个clientId会把对方踢下线。上线前必须用mosquitto_sub -i client-id与Broker的认证配置。EMQX默认允许匿名连接生产环境必须改成用户名密码认证并限制每个用户只能订阅/发布指定前缀的topic。我用EMQX的ACL规则做了按设备前缀的权限隔离防止设备端越权订阅其他设备的topic。第三断线重连参数调优。Paho的setAutomaticReconnect(true)只是打开了自动重连默认的重连间隔是1秒、2秒、4秒这样指数退避增长。但对于无人货柜这种设备我建议设置setMaxReconnectDelay(30000)把最大重连间隔控制在30秒内避免设备长时间无法恢复连接。第四消息积压监控。如果Broker上某个topic的消息生产速度快于消费速度消息就会堆积。我在SpringBoot里加了指标监控每次消费消息时通过Micrometer累加计数配合Prometheus和Grafana做告警。当消费速率低于生产速率超过1分钟时立刻触发告警防止消息越积越多导致系统雪崩。第五灰度发布策略。售货柜设备固件升级不能一把梭我先挑了两台测试柜升级新固件运行一周确认稳定后再推全量。服务端也是一样先切10%的设备流量到新版本代码观察订单成功率没有下降后再全量发布。4.5 一个值得记录的现场故障最后分享一个让我印象深刻的故障处理过程。某天下午运营反馈有3台柜机的所有扣款订单全部失败用户那边显示扣款成功但服务端账单显示支付失败。两边弹出来的金额都不一样。我第一反应是支付渠道回调出了问题于是查了支付网关日志发现回调根本没有到达服务端。进一步排查Nginx日志才发现那段时间服务器IP被机房判定为异常流量进行了限流所有HTTPS回调请求被丢弃了一部分。这个问题的根因不在代码但在架构上有明显隐患。之后我把支付回调接收做成了独立于主应用的轻量级服务走单独的Nginx转发路径配置独立的限流策略。这个教训告诉我做物联网支付系统不能只盯着MQTT这条链路传统Web链路的稳定性同样重要。写这篇文章的时候我脑子里过了一遍整个项目的点点滴滴。从最开始搭建SpringBoot工程、接入MQTT客户端到后面一点点补全设备管理、订单状态机、消息补偿机制每个环节都踩过不少坑。如果你也在用SpringBoot做物联网相关的项目我最想强调的一点是MQTT只是消息管道真正决定系统能否稳定运行的是管道两端的设计——设备端的容错能力和服务端的幂等补偿。把这两个点想透了这个项目就成功了一大半。