先交代下背景。这个项目不是一开始就成形的,最初只是帮一个做宠物用品的朋友调研“自助宠物洗浴机”到底有没有搞头,查了一圈发现市面上确实有类似共享洗车机的形态,但真正完整的、能跑通的公开项目源码很少。后来我干脆自己动手,从硬件选型到Java后端再到业务流程状态机,把一整套“无人共享宠物洗澡物联网系统”给落地出来了。本文就是把整个推导过程、技术方案、关键代码和踩坑记录整理出来,给想做物联网毕设、无人值守设备应用、或者对共享经济硬件落地感兴趣的朋友一个可以直接照着做的参考。
1. 无人宠物洗澡这个场景,凭什么需要一套物联网系统
1.1 共享洗浴面临的四大核心问题
先别急着聊技术,我们得先把场景想透。无人共享宠物洗澡机,你可以理解成“宠物版的共享洗车房”:一个封闭舱体,用户扫码付款后,机器自动放水、上沐浴露、循环冲洗、然后吹干。整个过程无人干预,但现实里这件事远比共享洗衣房麻烦得多。
- 过程不可控:宠物会动,水流方向和温度必须实时调节,出问题要能自动停机;
- 计费不可糊弄:按分钟计费还是按套餐计费?中途用户不满意退款怎么算?设备中途坏了怎么赔付?
- 设备会坏:水泵烧了、水管堵了、泡沫传感器误报,硬件异常需要第一时间上报并停止接单;
- 安全是底线:洗澡舱是封闭空间,断电、断网、宠物受惊失控,每一秒都得有状态兜底。
这四个问题拆开看,每一个都必须靠一套“懂业务”的软件系统来解决,而不是简单做几个继电器开关。所以这个项目本质上不是“给单片机写个流水灯控制程序”,而是一个完整的分布式物联网应用:设备端负责执行,云端负责决策,中间用通讯协议串联起来。
1.2 为什么后端服务要选Java而不是更“轻”的技术
说实话,物联网设备端最常见的组合是C/C++、MicroPython,后端则有人用Node-RED、Python Flask,甚至直接用云厂商的物联网平台托管。为什么这个项目偏偏强调“Java驱动”?
三个原因:
业务复杂度决定了必须有一个重型后端。计费、订单、用户余额、设备状态、异常事件、退款事务、渠道对账,这些全部是典型的企业级业务,Java生态在这一块的支撑最成熟。用脚本语言写原型很快,但等到要处理高并发扣费和事务一致性的时候,还是要回到Spring Boot这一套。
物联网毕业设计和面试场景高度匹配。当前热词里大量出现“java面试八股文”“java面试题”“物联网毕业设计”,说明大量读者其实是学生或刚转行的开发者。一个Java版的物联网后端,既能展示Spring Boot能力,又能展示Netty/MQTT/RocketMQ这类中间件知识,面试时有得讲。
硬件接入协议本身是语言无关的。设备端通过MQTT上报JSON数据,Java后端接收、解析、入库、触发业务动作,这一套流程非常标准。选Java不会带来额外成本,反而后续接支付、接短信通知、接管理后台时省事很多。
1.3 整套系统的数据链路长什么样
一句话概括:设备感知环境数据 → 通过MQTT上报到服务端 → Java后端做业务判断并落库 → 下发控制指令 → 设备执行动作 → 小程序实时展示状态。
除了这条主链路,还有一条辅助链路:设备心跳数据不停上报,服务端据此维护在线状态;同时服务端会把洗浴过程的温度、水位、用电量等数据同步给OneNET这类物联网可视化平台,方便远程监控和数据分析。
提示:这里把OneNET作为可视化监控层,而不是核心业务层,相当于给整套系统加了一块“仪表盘”,不会影响核心业务的独立性。如果你不想用OneNET,也可以换成Grafana+时序数据库的方案。
2. 硬件端的整体设计与控制核心:ESP32-S3主控和传感器布局
2.1 主控选型和外围传感器怎么搭配
设备端主控我选了ESP32-S3。为什么不是STM32?不是ESP8266?理由很直接:
- ESP32-S3支持Wi-Fi + BLE,可以直接联网,省掉额外的联网模块;
- 双核240MHz,跑状态机、处理传感器数据、维护MQTT长连接足够用;
- 外设接口丰富,控制电磁阀、水泵、风机都需要多路GPIO和PWM;
- 开发效率高,可以用Arduino框架也可以上ESP-IDF,社区资料多,学生团队上手快。
外围传感器和输出设备,我整理了一个清单,照着买基本不会错:
| 类别 | 具体器件 | 作用 | 注意要点 |
|---|---|---|---|
| 环境感知 | 防水NTC温度探头 | 实时监测水温 | 探头要直接接触水路,别只测空气温度 |
| 水位感知 | 压力式液位传感器 | 判断水箱水位和舱内积水 | 不要用浮球开关,沐浴露泡沫会导致误判 |
| 门控检测 | 磁簧开关 + 电磁锁 | 检测舱门开关状态和锁止 | 无人场景必须有锁止,防中途开门 |
| 动作执行 | 进水电磁阀、循环水泵、风机、雾化泵 | 控制进水、冲洗、吹干、沐浴露喷洒 | 水泵必须支持干转保护,否则容易烧毁 |
| 主控扩展 | 继电器模块 / MOSFET驱动板 | 放大GPIO驱动能力,控制大功率设备 | 水泵和风机不能直接接GPIO,必须隔离 |
| 可选 | 摄像头模块(OV2640) | 远程查看舱内状态 | 涉及隐私,建议默认关闭 |
核心控制逻辑是这样的:ESP32-S3开机后先完成自检,如果所有传感器状态正常,就连接Wi-Fi并建立MQTT长连接,然后进入待命状态。每次收到服务端下发的新指令(比如“开始进水”),控制器先校验当前设备状态是否符合执行条件,再驱动对应继电器吸合。
2.2 设备端状态机与“断网自救”逻辑
设备端最容易被忽视的就是断网自救。很多人的设计是“断网就停机等恢复”,这在无人值场景里是灾难:宠物洗到一半断网,风机停了,舱里积水泡着吹不干,用户投诉电话直接打爆。
我设计的设备端状态机包含以下几个状态:
- IDLE(空闲,可被预约)
- ACTIVATED(已激活,准备注水)
- WASHING(洗浴中,注水+喷洒沐浴露+循环冲洗)
- RINSING(漂洗中,清水循环)
- DRYING(吹干中,风机+加热)
- SELF_CLEAN(自清洁,舱体清洗)
- ERROR(故障停机,拒绝接单)
- OFFLINE(离线,网络断开)
关键点在于:洗浴过程中的任何时刻,如果MQTT断连,设备不会立即停机,而是维持当前状态继续完成当前阶段的剩余动作,并记录“断网事件”。比如正在DRYING阶段,断网了,风机继续吹够剩余时间,然后设备进入一个“待确认”状态,等网络恢复后把断网期间的状态补报给服务端。
提示:这个“断网继续执行”的决策要慎重。对宠物洗澡来说,维持环境安全是第一位的,风机、水泵这种安全设备可以继续运行;但注水阀门这种涉及水浸风险的,断网时绝不能继续开。所以我在状态机里额外加了一个规则:注水类动作必须在网络心跳正常时才允许执行。
2.3 设备和云端之间的MQTT协议设计
协议这块我直接采用MQTT,不用HTTP轮询。原因很现实:设备在洗浴过程中会高频上报数据(水温、液位、风机状态),如果每次都走HTTP握手,连接开销大,服务端压力也大,并且服务端想主动给设备下指令时,HTTP需要设备先发起请求才能响应,做不到真正的“实时推送”。
MQTT的主题设计如下:
| 主题 | 方向 | 消息内容 | 用途 |
|---|---|---|---|
| device/{deviceId}/telemetry | 设备→云 | 传感器数据JSON | 温度、液位、电源状态等 |
| device/{deviceId}/status | 设备→云 | 设备状态信息 | 当前状态机枚举值、错误码 |
| device/{deviceId}/heartbeat | 设备→云 | 心跳包 | 维持在线状态 |
| cloud/{deviceId}/command | 云→设备 | 控制指令 | 开始进水、启动风机、停机等 |
| cloud/{deviceId}/config | 云→设备 | 参数配置 | 修改洗浴时长、温度阈值 |
| device/{deviceId}/event | 设备→云 | 事件上报 | 断网、恢复、异常、自清洁完成 |
消息体统一使用JSON。为什么不用二进制协议?因为开发效率优先,JSON可读性好,调试方便,ESP32-S3上cJSON解析开销完全可以接受。当前业务量级下,完全没有必要为节省那几KB流量引入更复杂的编解码。
一个完整的上报消息示例:
{ "deviceId": "petwash_001", "ts": 1713500000000, "state": "WASHING", "waterTemp": 32.5, "waterLevel": 78, "pump": true, "fan": false, "doorLocked": true, "faultCode": 0 }3. Java后端的大脑中枢:设备接入、心跳保活与指令下发
3.1 Spring Boot + MQTT Broker的设备接入层
后端主框架我用Spring Boot 2.7,设备接入层用的是Eclipse Paho MQTT客户端连接EMQX Broker。为什么要单独搞一个Broker而不是让设备直连后端?因为MQTT本身就是发布/订阅模式,Broker负责维护海量长连接和主题分发,后端服务只管订阅关心的主题就行。这样设备量增长时,Broker可以水平扩展,后端业务服务不需要跟着改。
后端接入层的核心类是MqttMessageHandler,通过@MqttSubscribe注解订阅上面那张表里的所有主题:
@Component public class MqttMessageHandler { private static final Logger log = LoggerFactory.getLogger(MqttMessageHandler.class); @Resource private DeviceTelemetryService telemetryService; @Resource private DeviceStateMachineService stateMachineService; @Resource private DeviceCommandService commandService; @MqttSubscribe(topic = "device/+/telemetry") public void handleTelemetry(String topic, MqttMessage message) { String deviceId = TopicUtils.resolveDeviceId(topic); TelemetryData data = JsonUtils.parseObject(message.getPayload(), TelemetryData.class); // 上报数据先做基本校验,设备ID不存在直接丢弃 if (!deviceRegistry.exists(deviceId)) { log.warn("unknown device telemetry: {}", deviceId); return; } telemetryService.save(deviceId, data); // 温度超过45度或者水位低于安全阈值时,触发保护逻辑 if (data.getWaterTemp() > 45 || data.getWaterLevel() < 15) { commandService.sendCommand(deviceId, CommandType.EMERGENCY_STOP); } } @MqttSubscribe(topic = "device/+/status") public void handleStatus(String topic, MqttMessage message) { String deviceId = TopicUtils.resolveDeviceId(topic); DeviceStatus status = JsonUtils.parseObject(message.getPayload(), DeviceStatus.class); stateMachineService.processDeviceStatus(deviceId, status); } @MqttSubscribe(topic = "device/+/heartbeat") public void handleHeartbeat(String topic, MqttMessage message) { String deviceId = TopicUtils.resolveDeviceId(topic); deviceRegistry.refreshHeartbeat(deviceId); } @MqttSubscribe(topic = "device/+/event") public void handleEvent(String topic, MqttMessage message) { String deviceId = TopicUtils.resolveDeviceId(topic); DeviceEvent event = JsonUtils.parseObject(message.getPayload(), DeviceEvent.class); eventService.record(deviceId, event); } }处理思路很简单:收到遥测数据,先做基础校验,再落库,同时判断是否有需要紧急干预的异常数据;收到状态消息,就交给状态机服务区处理业务流转;心跳和事件则单独对接。
3.2 心跳、离线判定和服务端补偿
设备端每10秒发送一次心跳包,这个心跳不是空数据,里面带设备状态和最近一次遥测数据的关键摘要。服务端收到后刷新Redis中的设备心跳时间。
离线判定规则:超过30秒未收到心跳,设备标记为“疑似离线”;超过60秒未收到心跳,标记为“离线”。
为什么分成两档?因为Wi-Fi网络总会抖动,偶尔丢一两个包很正常。如果立刻判定离线并停止所有业务,用户在洗浴中途会莫名其妙被打断。所以“疑似离线”阶段只触发预警通知,不打断正在进行的正常业务;只有确认离线超过60秒,才会触发应急停机逻辑。
这里有一个隐藏的坑:设备断网后,服务端怎么判断它到底是“真的离线”还是“消息堆积没处理到”?解决方式是利用MQTT的遗嘱消息(Last Will and Testament)。设备在连接Broker时设置遗嘱主题device/{deviceId}/lwt,内容为offline;如果设备异常断网,Broker会立即替设备发布这条遗嘱消息,服务端订阅到后就可以第一时间确认设备离线,而不需要等60秒超时。
3.3 服务端下发指令的幂等设计
云端下发指令到设备,必须要考虑“消息丢失后重发会不会造成重复执行”。比如“开始进水”指令,如果发了两遍,设备就会执行两次,舱体会被淹。我这里的处理方式是:
- 每条指令带上全局唯一
commandId(UUID生成); - 设备端收到指令后先校验
commandId是否已执行过,如果重复直接忽略; - 服务端发送指令后把
commandId写入Redis,并设置30秒的过期时间;如果设备始终没有回复ACK,服务端会基于同样的commandId重发,不会生成新指令。
核心发送逻辑:
@Service public class DeviceCommandService { private static final Duration RETRY_EXPIRE = Duration.ofSeconds(30); @Resource private MqttGateway mqttGateway; @Resource private RedisTemplate<String, String> redisTemplate; public boolean sendCommand(String deviceId, CommandType type) { String commandId = UUID.randomUUID().toString(); return doSend(deviceId, commandId, type); } private boolean doSend(String deviceId, String commandId, CommandType type) { CommandMessage command = new CommandMessage(); command.setCommandId(commandId); command.setType(type.getCode()); command.setTimestamp(System.currentTimeMillis()); // 发送前先登记,防止重复下发 String dedupKey = "cmd:dedup:" + deviceId + ":" + commandId; Boolean first = redisTemplate.opsForValue() .setIfAbsent(dedupKey, "1", RETRY_EXPIRE); if (Boolean.FALSE.equals(first)) { log.warn("duplicate command ignored: {}", commandId); return true; } mqttGateway.publish("cloud/" + deviceId + "/command", JsonUtils.toJsonString(command)); return true; } }这个设计应对的是最常见的网络抖动场景,保证“至少一次送达,且最多一次执行”。
4. 洗浴全流程状态机:无人值守的关键就是把状态拆到最细
4.1 六个业务状态与转换条件
无人共享宠物洗澡系统的业务状态,不能等于硬件状态。硬件状态描述的是设备当前在干什么,而业务状态描述的是“用户这一单当前进行到哪一步了”。前者是物理事实,后者是商业逻辑。
我设计了一套独立的业务状态机:
| 业务状态 | 含义 | 触发条件 | 下一步 |
|---|---|---|---|
| WAITING_START | 用户已扫码下单,等待设备就绪 | 支付成功回调 | 下发启动指令 |
| PREPARING | 设备预热、注水、检测环境 | 收到启动ACK | 进入洗浴中 |
| BATHING | 洗浴进行中(可细分阶段) | 预热完成 | 进入吹干/结算 |
| SUSPENDED | 暂停(用户要求/异常中断) | 暂停指令/异常 | 恢复/终止 |
| COMPLETED | 订单完成,等待设备自清洁 | 流程正常结束 | 触发自清洁 |
| REFUNDED | 已退款关闭 | 用户取消/客服介入 | 终态 |
为什么服务端要自己再维护一套业务状态,而不是直接用硬件上报的状态?
因为服务端必须对每一单负责。想象一个场景:用户在手机上点了“暂停”,服务端必须立刻记下当前业务状态,通知设备暂停水泵和风机,并开始计时“暂停时长”。如果这个判断完全交给设备端做,一旦设备状态机出错,服务端根本不知道,就会出现用户被多扣费或者少扣费的问题。
服务端状态机核心流转代码:
public class WashOrderStateMachine { private final Map<OrderStatus, Set<OrderStatus>> transitions = new EnumMap<>(OrderStatus.class); public WashOrderStateMachine() { // 定义允许的状态流转路径 transitions.put(OrderStatus.WAITING_START, EnumSet.of(OrderStatus.PREPARING, OrderStatus.REFUNDED)); transitions.put(OrderStatus.PREPARING, EnumSet.of(OrderStatus.BATHING, OrderStatus.SUSPENDED, OrderStatus.REFUNDED)); transitions.put(OrderStatus.BATHING, EnumSet.of(OrderStatus.SUSPENDED, OrderStatus.COMPLETED, OrderStatus.REFUNDED)); transitions.put(OrderStatus.SUSPENDED, EnumSet.of(OrderStatus.BATHING, OrderStatus.REFUNDED)); transitions.put(OrderStatus.COMPLETED, EnumSet.of(OrderStatus.PREPARING)); // 自清洁后可以接下一单 } public synchronized OrderStatus transition(OrderEntity order, OrderStatus target) { Set<OrderStatus> allowed = transitions.get(order.getStatus()); if (allowed == null || !allowed.contains(target)) { throw new IllegalStateException( String.format("非法流转: %s -> %s", order.getStatus(), target)); } order.setStatus(target); return target; } }状态机的严谨程度直接决定业务安全性。非法流转直接抛异常,宁可让订单卡住等人工介入,也不能让它走到一个语义混乱的状态。
4.2 中途退款、异常中断、余额不足怎么处理
无人设备最怕的就是钱算不明白。我在设计退款和异常处理时,遵循了几个原则:
先暂停再退款、先检测再赔付。用户申请中途退款,不直接退全款,而是先把设备切换到SUSPENDED状态,停止计费;服务端此时读取设备上报的水温、水位、风机运行时长、总体运行时间,根据这些数据计算“已产生费用”,退还剩余部分。比如用户洗了15分钟,套餐是30元包60分钟,但用户因宠物不配合申请终止,系统只扣除前15分钟对应的7.5元,剩余22.5元原路退回。
计费单位按分钟取整,不足一分钟按一分钟计。服务端每分钟执行一次扣费任务,读取订单的开始时间,计算已用分钟数,和上次扣费的分钟数相减,差额即为本次扣费金额。这样实现简单,且误差在可接受范围。
余额不足时不直接停机。宠物正在吹干,余额不足如果直接停机,宠物会受凉。所以系统会先发送“余额不足提醒”,同时把业务状态切到SUSPENDED;此时风机继续运转,但水循环和加热停止,给用户3分钟缓冲时间决定是否续费。如果3分钟后仍未续费,设备才会执行“安全停机”,并通知用户“订单未完成,请现场处理”。
4.3 监控看板与OneNET数据可视化
出于设备运维的考虑,我把遥测数据额外同步到OneNET平台,用来做趋势可视化和历史数据查询。主要看三个指标:水温曲线、水位曲线、风机电流曲线。
为什么要做这一步?因为设备故障往往不是突然发生的,而是有迹可循的。水温持续偏低,可能是加热管老化;水位下降速度异常,可能是排水阀没关严;风机电流偏高,可能是进风口堵塞。这些趋势在OneNET的折线图上一眼就能看出来,比等设备报故障码再处理高效得多。
同步逻辑很简单:后端收到遥测数据后,通过OneNET提供的API将数据点写入对应数据流,不需要引入额外SDK,用RestTemplate就行:
public void syncToOneNET(String deviceId, TelemetryData data) { OneNetDataPoint point = new OneNetDataPoint(); point.setDatastreams(List.of( new DataStream("water_temp", data.getWaterTemp()), new DataStream("water_level", data.getWaterLevel()), new DataStream("fan_current", data.getFanCurrent()) )); String url = "http://api.heclouds.com/devices/" + deviceId + "/datapoints"; restTemplate.postForEntity(url, point, String.class); }实际情况中这台设备的遥测数据不一定要走OneNET做业务,但多一条数据通道总是好事,特别是当你需要给投资人、导师或者运维团队展示“这套系统是真的在运行”的时候,一条漂亮的实时曲线比一百页PPT都管用。
5. 计费、自清洁与多设备并发:真正量产才会踩到的坑
5.1 预授权冻结+按分钟计费的实现逻辑
计费方案我直接对标共享充电宝和共享洗衣机的做法:用户扫码后先冻结套餐金额(比如30元),洗浴结束后按实际时长从冻结金额中扣除,剩余部分自动解冻。
这样做的好处是:用户不需要在洗到一半时再进行一次支付操作,体验顺畅;对商家来说,资金有保障,不会出现用户洗了100分钟只付了30元就跑了的情况。
预授权走微信支付的“押金冻结/解冻”接口。订单创建时调用冻结接口,订单完成时调用解冻接口,并同步发起一笔实际扣款。这个流程里最需要注意的是回调顺序问题:
注意:微信支付的回调是异步的,冻结回调、解冻回调、扣款回调三者到达服务端的顺序不保证。我踩过一次坑,扣款回调先到,解冻回调后到,结果是用户被多扣了钱。解决方案是:统一以订单业务状态为准,只有COMPLETED状态的订单才允许执行扣款和解冻;其余回调一律挂起等待状态机流转到COMPLETED再处理。
按分钟扣费的核心代码:
@Component public class BillingTask { @Resource private WashOrderMapper orderMapper; @Resource private BalanceService balanceService; @Scheduled(fixedRate = 60000) public void chargePerMinute() { List<WashOrder> activeOrders = orderMapper.selectActiveOrders(); for (WashOrder order : activeOrders) { long usedMinutes = Duration.between(order.getStartTime(), LocalDateTime.now()).toMinutes() + 1; long lastChargedMinutes = order.getChargedMinutes(); long diff = usedMinutes - lastChargedMinutes; if (diff <= 0) { continue; } BigDecimal amount = order.getUnitPricePerMinute().multiply(BigDecimal.valueOf(diff)); boolean success = balanceService.deduct(order.getUserId(), amount, "洗浴计费", order.getOrderId()); if (success) { order.setChargedMinutes(usedMinutes); orderMapper.updateById(order); } else { // 余额不足,进入预警流程 notifyService.notifyLowBalance(order); } } } }5.2 自清洁时序与硬件保护
订单完成后,设备并不会立刻进入空闲状态,而是先执行自清洁:用清水冲洗舱内地板和循环管道,再开启风机吹干残留水珠,整个流程约3分钟。自清洁期间设备不接受新订单,状态显示为“维护中”。
自清洁的时序控制放在设备端,服务端只发一条“启动自清洁”指令,具体控制逻辑由ESP32-S3本地执行:
- 关闭排水阀,注水至低液位;
- 开启循环泵,清水循环冲刷管道2分钟;
- 打开排水阀,排空污水;
- 关闭排水阀,二次注水至低液位,再次循环1分钟;
- 排空;
- 打开风机,吹干残留水分30秒;
- 自清洁完成,上报状态,设备回到IDLE。
这个时序里最容易出的bug是:设备端自清洁时序执行到一半,服务端误以为设备已经空闲,直接派了下一单。解决方式是在服务端增加一个“设备忙碌”锁,自清洁期间设备状态为BUSY,任何订单预约都会返回“设备维护中”,等自清洁完成事件上报后才解除。
5.3 多设备并发和消息乱序的应对
当你同时接入10台、20台设备时,消息并发和顺序问题就会浮出水面。最初我以为每台设备独立上报,互不干扰,结果上线测试时发现一个诡异问题:同一台设备的状态消息,上报顺序是WASHING → RINSING → WASHING,服务端收到RINSING后把订单状态切到了BATHING的漂洗子状态,随后又收到一条延迟的WASHING,状态机不认这个非法流转,直接抛异常,导致订单卡死。
这个问题的根源是:MQTT不保证同一设备的消息在Broker到服务端的消费链路上严格有序。设备端发了三条消息,可能因为网络抖动原因,第三条先被服务端消费,第二条后消费。
解决办法有三层:
- 设备端给每条消息带上自增序号
seq,服务端对同一设备的消息做seq去重和乱序缓存; - 服务端状态机比较消息序号,只处理比当前已处理序号更新的消息;
- 如果检测到序号跳跃过大,主动向设备发起“状态同步”请求,让设备重新上报当前状态。
最终实现上,我在Redis里维护每个设备的“最近已处理消息序号”,收到新消息后先判断:
public boolean isStaleMessage(String deviceId, long seq) { String key = "device:seq:" + deviceId; Long lastSeq = redisTemplate.opsForValue().increment(key); return seq <= lastSeq; }这种乱序处理是物联网项目里非常典型的问题,处理不好就是线上事故。我在这里卡了整整两天,最后靠给消息加单调递增序号才彻底解决。
6. 源码项目结构说明与复现这个项目的建议
6.1 后端工程目录怎么划分
整个后端工程基于Maven构建,按模块化方式拆分,核心目录结构如下:
pet-wash-iot-backend ├── pom.xml ├── src/main/java/com/petwash/iot │ ├── common // 通用工具、常量、异常处理 │ ├── config // Spring配置、MQTT配置、Redis配置 │ ├── controller // 面向小程序/管理后台的REST接口 │ ├── mqtt // MQTT消息处理、主题解析 │ ├── device // 设备注册、心跳管理、指令下发 │ ├── order // 订单服务、业务状态机 │ ├── billing // 计费、退款、预授权 │ └── monitor // 遥测数据、OneNET同步、告警规则 ├── src/main/resources │ ├── mapper // MyBatis映射文件 │ └── application.yml └── sql // 建表脚本这个划分的核心思路是按业务域划分,而不是按技术层划分。早期我把类按controller/service/mapper这样堆,结果订单、设备、计费三个业务域的逻辑耦合在一起,改一个需求牵一发动全身。按业务域拆分后,每块代码自己能独立演进。
设备接入层的核心链路是:EMQX Broker收到设备消息 → 根据主题前缀分发到对应处理器 → 消息解析成Java对象 → 写入数据库/Redis → 触发业务状态机。整条链路在mqtt包和device包里完成。
6.2 核心代码片段:设备上报消息处理
在3.1节的MqttMessageHandler基础上,再补充服务端收到设备状态上报后,如何驱动业务状态机流转的关键代码。重点是区分“设备硬件状态”和“服务端业务状态”,前者用DeviceStateMachineService处理,后者用WashOrderStateMachine处理:
@Service public class DeviceStateMachineService { @Resource private WashOrderMapper orderMapper; @Resource private DeviceCommandService commandService; public void processDeviceStatus(String deviceId, DeviceStatus status) { WashOrder activeOrder = orderMapper.selectActiveByDeviceId(deviceId); if (activeOrder == null) { // 没有进行中的订单,收到设备状态只记录,不处理 log.debug("no active order for device {}", deviceId); return; } switch (activeOrder.getStatus()) { case PREPARING: // 设备进入WASHING说明预热完成,业务状态推进到BATHING if (status.getState() == DeviceState.WASHING) { orderMapper.updateStatus(activeOrder.getOrderId(), OrderStatus.PREPARING, OrderStatus.BATHING); } break; case BATHING: // 设备完成DRYING,说明洗浴流程结束,业务状态推进到COMPLETED if (status.getState() == DeviceState.IDLE && activeOrder.getPaidAmount() != null) { orderMapper.updateStatus(activeOrder.getOrderId(), OrderStatus.BATHING, OrderStatus.COMPLETED); // 触发计费结算和自清洁 billingService.settle(activeOrder); commandService.sendCommand(deviceId, CommandType.START_SELF_CLEAN); } break; case SUSPENDED: // 处于暂停状态的订单,收到设备的打断恢复信号时自动恢复 if (status.getState() == DeviceState.WASHING) { orderMapper.updateStatus(activeOrder.getOrderId(), OrderStatus.SUSPENDED, OrderStatus.BATHING); } break; default: break; } } }这个类承担的职责是桥接硬件状态和业务状态,也是整个系统里最容易出错的地方。因为硬件上报的状态是物理事实,而业务状态是契约逻辑,两者之间存在一个隐式的对应关系。设计时必须把每个对应关系写清楚,不能靠“感觉应该差不多”。
6.3 复现这个项目时,硬件和代码的衔接要点
最后给想动手复现的朋友几个实操建议,都是我自己踩过的坑:
先跑通“模拟器”,再碰真硬件。硬件调试非常耗费时间,而且新手很容易把传感器接反烧掉。我的做法是先用一个Java写个Device Simulator,模拟MQTT消息上报和指令接收,把后端业务逻辑全部调通之后,再让ESP32-S3对接真硬件。这个顺序能节省至少一周时间。
ESP32-S3的Wi-Fi连接稳定性需要专门处理。设备如果放在铁皮舱体内,Wi-Fi信号衰减极其严重。我建议在舱体顶部外置天线,或者给ESP32-S3接一个外置天线底座。同时,固件里要做Wi-Fi断线自动重连,并且把重连逻辑单独跑在一个独立任务里,不能让主循环卡在Wi-Fi连接上。
一定要有RTC和本地事件日志。设备每次上报都要带本地时间戳,同时把状态切换记录在SD卡或SPI Flash里。不然等出问题回溯的时候,设备本地日志和服务端日志完全对不上,排查非常痛苦。
谨慎使用OTA自动升级。无人设备的固件升级一旦失败,设备变砖,需要有人现场恢复。我建议量产阶段不要全量推送自动升级,先灰度一台,确认稳定后再逐步扩大。
接口设计要预留扩展位。比如以后可能增加“宠物自助称重”功能,或者接入摄像头做“异常行为识别”,那MQTT消息体里的字段就要预留可扩展的Map结构,不要写死字段名。
做完这一整套项目,我最大的感受是:无人共享宠物洗澡看起来只是“把人工做的事自动化”,但真正难的不是自动化,而是“出了问题没人管的时候,系统怎么保证不伤宠物、不伤设备、不糊弄账”。这三个“不”字,就是这套Java物联网系统存在的全部意义。源码本身的每一行代码,其实都是在为一件事服务——让机器在无人值守的状态下,自己知道该做什么、不该做什么。