SpringCloud Alibaba无人售货柜实战(五):设备通信协议设计——MQTT/HTTP指令下发与状态回调
让售货柜开门它就开门,让它重启它就重启——这背后需要一套严谨的通信协议。指令丢了怎么办?设备没响应怎么办?这篇全给你兜住。
一、设备通信架构
整个通信链路是一条完整的"指令生命周期":
服务端下发指令 │ ▼ MQTT Broker → device/{sn}/command Topic │ ▼ 设备端接收 → 执行操作(开电磁锁/重启等) │ ▼ 设备端上报回调 → device/{sn}/callback Topic │ ▼ 服务端处理回调 → 更新指令状态 → 触发后续业务正常情况下这条链路在2秒内跑完。但现实世界有网络抖动、设备死机、MQTT断连等各种意外,所以通信协议必须设计好超时、重试、幂等三道保险。
二、通信协议设计原则
- 简洁:字段名短小精悍,JSON层级不超过3层,减少设备端解析负担
- 可靠:每条指令有唯一ID,支持幂等执行和结果追踪
- 可扩展:预留
extra字段,新增指令类型不改协议结构 - 可追踪:每条指令从下发到回调全链路有日志,方便排查
三、下行指令协议
3.1 指令结构定义
服务端发给设备的指令格式:
{"commandId":"cmd-550e8400-e29b-41d4-a716-446655440000","command":"OPEN_DOOR","params":{"orderId":"202607291234567890","maxDuration":300},"timeout":30,"timestamp":1753766400000,"sign":"a1b2c3d4e5f6"}| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
| commandId | String | 是 | 指令唯一ID,UUID生成,用于关联回调 |
| command | String | 是 | 指令类型枚举 |
| params | Object | 否 | 指令参数,不同指令参数不同 |
| timeout | int | 是 | 超时时间(秒),默认30 |
| timestamp | long | 是 | 下发时间戳,设备端可用于防重放 |
| sign | String | 是 | 签名,MD5(commandId + command + timestamp + secret) |
3.2 指令类型定义
| 指令类型 | 说明 | params参数 | 超时建议 |
|---|---|---|---|
OPEN_DOOR | 开柜门 | orderId(订单号), maxDuration(最大开门时长秒) | 10秒 |
CLOSE_DOOR | 强制关柜门 | 无 | 10秒 |
RESTART | 重启设备 | delay(延迟秒数) | 60秒 |
SYNC_TIME | 同步时间 | serverTime(服务器时间戳) | 5秒 |
INVENTORY | 盘点指令 | 无(设备端返回当前库存) | 30秒 |
UPDATE_CONFIG | 更新配置 | heartbeatInterval, volume, autoClose… | 10秒 |
UPLOAD_LOG | 上传日志 | startTime, endTime | 60秒 |
TAKE_PHOTO | 拍照 | cameraId(摄像头编号) | 10秒 |
四、上行回调协议
设备执行完指令后,通过回调Topic上报执行结果:
{"commandId":"cmd-550e8400-e29b-41d4-a716-446655440000","status":"SUCCESS","data":{"doorOpen":true,"openDuration":45},"errorCode":null,"errorMsg":null,"timestamp":1753766402000}| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
| commandId | String | 是 | 关联的指令ID,和下行指令一一对应 |
| status | String | 是 | SUCCESS / FAILED / TIMEOUT / UNSUPPORTED |
| data | Object | 否 | 执行结果数据,不同指令返回不同 |
| errorCode | String | 否 | 失败时的错误码 |
| errorMsg | String | 否 | 失败时的错误描述 |
| timestamp | long | 是 | 回调时间戳 |
4.1 各指令的回调data定义
| 指令 | 回调data |
|---|---|
| OPEN_DOOR | {"doorOpen": true, "openDuration": 45} |
| CLOSE_DOOR | {"doorClosed": true} |
| RESTART | {"restartScheduled": true} |
| SYNC_TIME | {"synced": true, "deviceTime": 1753766402000} |
| INVENTORY | {"items": [{"productId": "P001", "count": 5}, ...]} |
| TAKE_PHOTO | {"imageUrl": "http://minio.xxx/photo/cmd-xxx.jpg"} |
五、指令下发Service
@Slf4j@ServicepublicclassDeviceCommandService{@AutowiredprivateDeviceCommandMappercommandMapper;@AutowiredprivateMqttGatewaymqttGateway;@AutowiredprivateRedisUtilsredisUtils;privatestaticfinalStringCOMMAND_PENDING_PREFIX="cmd:pending:";privatestaticfinalStringDEVICE_TOKEN_PREFIX="device:token:";/** * 下发指令 */publicDeviceCommandsendCommand(Stringsn,Stringcommand,JSONObjectparams,inttimeout){// 1. 生成指令IDStringcommandId="cmd-"+UUID.randomUUID().toString();// 2. 签名Stringtoken=redisUtils.get(DEVICE_TOKEN_PREFIX+sn);Stringsign=SecureUtil.md5(commandId+command+System.currentTimeMillis()+token);// 3. 构建指令消息JSONObjectmessage=newJSONObject();message.put("commandId",commandId);message.put("command",command);message.put("params",params);message.put("timeout",timeout);message.put("timestamp",System.currentTimeMillis());message.put("sign",sign);// 4. 存入数据库DeviceCommandcmd=newDeviceCommand();cmd.setCommandId(commandId);cmd.setDeviceSn(sn);cmd.setCommand(command);cmd.setParams(params.toJSONString());cmd.setStatus(0);// 待执行cmd.setTimeoutSeconds(timeout);cmd.setSendTime(LocalDateTime.now());commandMapper.insert(cmd);// 5. 通过MQTT下发Stringtopic="device/"+sn+"/command";mqttGateway.sendToMqtt(topic,message.toJSONString());log.info("指令已下发: sn={}, commandId={}, command={}",sn,commandId,command);// 6. 存入Redis待回调集合(用于超时检查)redisUtils.set(COMMAND_PENDING_PREFIX+commandId,sn,timeout+10,TimeUnit.SECONDS);// 7. 更新指令状态为已下发cmd.setStatus(1);commandMapper.updateById(cmd);returncmd;}/** * 发送开门指令(业务封装) */publicDeviceCommandopenDoor(Stringsn,StringorderId){JSONObjectparams=newJSONObject();params.put("orderId",orderId);params.put("maxDuration",300);returnsendCommand(sn,"OPEN_DOOR",params,10);}}六、回调处理
@Slf4j@ServicepublicclassDeviceCallbackService{@AutowiredprivateDeviceCommandMappercommandMapper;@AutowiredprivateRedisUtilsredisUtils;@AutowiredprivateOrderFeignClientorderFeignClient;privatestaticfinalStringCOMMAND_PENDING_PREFIX="cmd:pending:";/** * 监听设备回调 */@MqttMessageListener(topic="device/+/callback")publicvoidonCallback(MqttMessagemessage){Stringtopic=message.getTopic();Stringsn=topic.split("/")[1];Stringpayload=newString(message.getPayload(),StandardCharsets.UTF_8);CallbackReqreq=JSON.parseObject(payload,CallbackReq.class);log.info("收到设备回调: sn={}, commandId={}, status={}",sn,req.getCommandId(),req.getStatus());// 1. 查询指令DeviceCommandcmd=commandMapper.selectByCommandId(req.getCommandId());if(cmd==null){log.error("回调指令不存在: commandId={}",req.getCommandId());return;}// 2. 幂等检查:已经处理过的回调直接忽略if(cmd.getStatus()==2||cmd.getStatus()==3){log.warn("指令已处理,忽略重复回调: commandId={}, status={}",req.getCommandId(),cmd.getStatus());return;}// 3. 更新指令状态if("SUCCESS".equals(req.getStatus())){cmd.setStatus(2);// 成功}else{cmd.setStatus(3);// 失败}cmd.setResultData(req.getData()!=null?req.getData().toJSONString():null);cmd.setCallbackTime(LocalDateTime.now());commandMapper.updateById(cmd);// 4. 清除Redis待回调标记redisUtils.delete(COMMAND_PENDING_PREFIX+req.getCommandId());// 5. 触发后续业务handleCommandResult(sn,cmd,req);}/** * 根据指令类型触发后续业务 */privatevoidhandleCommandResult(Stringsn,DeviceCommandcmd,CallbackReqreq){switch(cmd.getCommand()){case"OPEN_DOOR":if("SUCCESS".equals(req.getStatus())){// 开门成功,通知订单服务orderFeignClient.onDoorOpened(cmd.getParamsObject().getString("orderId"));}else{// 开门失败,通知订单服务取消订单orderFeignClient.onDoorOpenFailed(cmd.getParamsObject().getString("orderId"),req.getErrorMsg());}break;case"INVENTORY":// 盘点结果同步到库存服务break;case"RESTART":log.info("设备重启指令已确认: sn={}",sn);break;}}}七、指令超时处理
指令下发后不是万事大吉——设备可能没收到、可能收到了但执行卡死了。必须有超时检查机制。
7.1 延迟队列方案
用RocketMQ的延迟消息实现超时检查:
@Slf4j@ServicepublicclassCommandTimeoutChecker{@AutowiredprivateDeviceCommandMappercommandMapper;@AutowiredprivateRocketMQTemplaterocketMQTemplate;@AutowiredprivateDeviceCommandServicecommandService;privatestaticfinalStringTIMEOUT_TOPIC="command-timeout-check";privatestaticfinalintMAX_RETRY=2;/** * 下发指令时发送延迟消息(延迟时间=指令超时时间) */publicvoidsendTimeoutCheck(StringcommandId,intdelaySeconds){Message<String>msg=MessageBuilder.withPayload(commandId).build();// RocketMQ延迟级别: 1s=1, 5s=2, 10s=3, 30s=4, 1m=5...intdelayLevel=delaySeconds<=5?2:(delaySeconds<=10?3:4);rocketMQTemplate.asyncSend(TIMEOUT_TOPIC,msg,newSendCallback(){@OverridepublicvoidonSuccess(SendResultsendResult){}@OverridepublicvoidonException(Throwablee){log.error("超时检查消息发送失败: commandId={}",commandId,e);}},3000,delayLevel);}/** * 消费超时检查消息 */@RocketMQMessageListener(topic=TIMEOUT_TOPIC,consumerGroup="command-timeout-group")@ComponentpublicclassTimeoutConsumerimplementsRocketMQListener<String>{@OverridepublicvoidonMessage(StringcommandId){DeviceCommandcmd=commandMapper.selectByCommandId(commandId);if(cmd==null)return;// 指令已完成(成功或失败),无需处理if(cmd.getStatus()==2||cmd.getStatus()==3){return;}log.warn("指令超时未回调: commandId={}, command={}, retryCount={}",commandId,cmd.getCommand(),cmd.getRetryCount());if(cmd.getRetryCount()<MAX_RETRY){// 重试:重新下发指令cmd.setRetryCount(cmd.getRetryCount()+1);cmd.setStatus(1);commandMapper.updateById(cmd);// 重新通过MQTT下发JSONObjectmessage=buildCommandMessage(cmd);mqttGateway.sendToMqtt("device/"+cmd.getDeviceSn()+"/command",message.toJSONString());// 再次发送延迟检查sendTimeoutCheck(commandId,cmd.getTimeoutSeconds());}else{// 超过最大重试次数,标记超时cmd.setStatus(4);// 超时commandMapper.updateById(cmd);log.error("指令最终超时: commandId={}",commandId);// 通知业务方处理}}}}八、HTTP备选通道
MQTT不可用时(Broker挂了或网络断了),设备通过HTTP轮询兜底拉取指令。
8.1 设备端轮询逻辑
设备端如果MQTT连接失败,自动降级为HTTP轮询模式:
每10秒请求: GET /api/device/{sn}/commands/pending 拉取待执行指令 → 执行 → POST /api/device/{sn}/callback 上报结果8.2 服务端轮询接口
@RestController@RequestMapping("/api/device")publicclassDevicePollController{@AutowiredprivateDeviceCommandMappercommandMapper;/** * 设备拉取待执行指令 */@GetMapping("/{sn}/commands/pending")publicResult<List<DeviceCommand>>getPendingCommands(@PathVariableStringsn){// 查询状态为"已下发"且未回调的指令List<DeviceCommand>commands=commandMapper.selectList(newLambdaQueryWrapper<DeviceCommand>().eq(DeviceCommand::getDeviceSn,sn).eq(DeviceCommand::getStatus,1).orderByAsc(DeviceCommand::getSendTime).last("LIMIT 5"));returnResult.success(commands);}/** * 设备HTTP上报回调 */@PostMapping("/{sn}/callback")publicResult<Void>callback(@PathVariableStringsn,@RequestBodyCallbackReqreq){callbackService.onCallback(sn,req);returnResult.success();}}HTTP轮询是兜底方案,不是常态。MQTT恢复后设备自动切回MQTT模式。双通道设计保证了通信可靠性。
九、安全设计
9.1 设备Token认证
设备连接MQTT时用Token做密码认证。EMQX配置用户认证后端,对接Redis验证:
MQTT连接用户名: {设备SN} MQTT连接密码: {Token} EMQX认证逻辑: GET device:token:{sn} → 比对密码9.2 指令签名防伪造
每条指令带sign字段,设备端验签后才执行:
// 设备端验签(Android/Java伪代码)publicbooleanverifySign(JSONObjectcommand,Stringtoken){StringcommandId=command.getString("commandId");Stringcmd=command.getString("command");longtimestamp=command.getLong("timestamp");Stringsign=command.getString("sign");StringexpectedSign=MD5Utils.md5(commandId+cmd+timestamp+token);returnexpectedSign.equals(sign);}9.3 防重放攻击
设备端维护一个最近100条commandId的LRU缓存,收到重复commandId直接忽略。配合timestamp字段,超过5分钟的指令直接丢弃。
十、通信协议完整定义表
| 指令 | 方向 | params | 回调data | 超时 | 重试 |
|---|---|---|---|---|---|
| OPEN_DOOR | 下行 | orderId, maxDuration | doorOpen, openDuration | 10s | 2次 |
| CLOSE_DOOR | 下行 | 无 | doorClosed | 10s | 1次 |
| RESTART | 下行 | delay | restartScheduled | 60s | 0次 |
| SYNC_TIME | 下行 | serverTime | synced, deviceTime | 5s | 1次 |
| INVENTORY | 下行 | 无 | items[] | 30s | 1次 |
| UPDATE_CONFIG | 下行 | 多个配置项 | updated | 10s | 1次 |
| UPLOAD_LOG | 下行 | startTime, endTime | logUrl | 60s | 0次 |
| TAKE_PHOTO | 下行 | cameraId | imageUrl | 10s | 1次 |
十一、小结
设备通信协议设计的核心就四个字:可靠、幂等。commandId贯穿整个生命周期,从下发到回调到超时检查,全靠它串联。MQTT是主通道,HTTP轮询是兜底,RocketMQ延迟消息做超时检查,三层保障确保指令不丢、不重、不卡。安全层面Token认证+指令签名+防重放三管齐下。这套协议跑通了,设备端和服务端就能稳定对话,后面的业务逻辑就是水到渠成的事。