SpringCloud Alibaba无人售货柜实战(五):设备通信协议设计——MQTT/HTTP指令下发与状态回调

发布时间:2026/8/2 0:33:25

SpringCloud Alibaba无人售货柜实战(五):设备通信协议设计——MQTT/HTTP指令下发与状态回调 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}字段类型必填说明commandIdString是指令唯一IDUUID生成用于关联回调commandString是指令类型枚举paramsObject否指令参数不同指令参数不同timeoutint是超时时间秒默认30timestamplong是下发时间戳设备端可用于防重放signString是签名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, endTime60秒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}字段类型必填说明commandIdString是关联的指令ID和下行指令一一对应statusString是SUCCESS / FAILED / TIMEOUT / UNSUPPORTEDdataObject否执行结果数据不同指令返回不同errorCodeString否失败时的错误码errorMsgString否失败时的错误描述timestamplong是回调时间戳4.1 各指令的回调data定义指令回调dataOPEN_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}五、指令下发ServiceSlf4jServicepublicclassDeviceCommandService{AutowiredprivateDeviceCommandMappercommandMapper;AutowiredprivateMqttGatewaymqttGateway;AutowiredprivateRedisUtilsredisUtils;privatestaticfinalStringCOMMAND_PENDING_PREFIXcmd:pending:;privatestaticfinalStringDEVICE_TOKEN_PREFIXdevice:token:;/** * 下发指令 */publicDeviceCommandsendCommand(Stringsn,Stringcommand,JSONObjectparams,inttimeout){// 1. 生成指令IDStringcommandIdcmd-UUID.randomUUID().toString();// 2. 签名StringtokenredisUtils.get(DEVICE_TOKEN_PREFIXsn);StringsignSecureUtil.md5(commandIdcommandSystem.currentTimeMillis()token);// 3. 构建指令消息JSONObjectmessagenewJSONObject();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. 存入数据库DeviceCommandcmdnewDeviceCommand();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下发Stringtopicdevice/sn/command;mqttGateway.sendToMqtt(topic,message.toJSONString());log.info(指令已下发: sn{}, commandId{}, command{},sn,commandId,command);// 6. 存入Redis待回调集合用于超时检查redisUtils.set(COMMAND_PENDING_PREFIXcommandId,sn,timeout10,TimeUnit.SECONDS);// 7. 更新指令状态为已下发cmd.setStatus(1);commandMapper.updateById(cmd);returncmd;}/** * 发送开门指令业务封装 */publicDeviceCommandopenDoor(Stringsn,StringorderId){JSONObjectparamsnewJSONObject();params.put(orderId,orderId);params.put(maxDuration,300);returnsendCommand(sn,OPEN_DOOR,params,10);}}六、回调处理Slf4jServicepublicclassDeviceCallbackService{AutowiredprivateDeviceCommandMappercommandMapper;AutowiredprivateRedisUtilsredisUtils;AutowiredprivateOrderFeignClientorderFeignClient;privatestaticfinalStringCOMMAND_PENDING_PREFIXcmd:pending:;/** * 监听设备回调 */MqttMessageListener(topicdevice//callback)publicvoidonCallback(MqttMessagemessage){Stringtopicmessage.getTopic();Stringsntopic.split(/)[1];StringpayloadnewString(message.getPayload(),StandardCharsets.UTF_8);CallbackReqreqJSON.parseObject(payload,CallbackReq.class);log.info(收到设备回调: sn{}, commandId{}, status{},sn,req.getCommandId(),req.getStatus());// 1. 查询指令DeviceCommandcmdcommandMapper.selectByCommandId(req.getCommandId());if(cmdnull){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_PREFIXreq.getCommandId());// 5. 触发后续业务handleCommandResult(sn,cmd,req);}/** * 根据指令类型触发后续业务 */privatevoidhandleCommandResult(Stringsn,DeviceCommandcmd,CallbackReqreq){switch(cmd.getCommand()){caseOPEN_DOOR:if(SUCCESS.equals(req.getStatus())){// 开门成功通知订单服务orderFeignClient.onDoorOpened(cmd.getParamsObject().getString(orderId));}else{// 开门失败通知订单服务取消订单orderFeignClient.onDoorOpenFailed(cmd.getParamsObject().getString(orderId),req.getErrorMsg());}break;caseINVENTORY:// 盘点结果同步到库存服务break;caseRESTART:log.info(设备重启指令已确认: sn{},sn);break;}}}七、指令超时处理指令下发后不是万事大吉——设备可能没收到、可能收到了但执行卡死了。必须有超时检查机制。7.1 延迟队列方案用RocketMQ的延迟消息实现超时检查Slf4jServicepublicclassCommandTimeoutChecker{AutowiredprivateDeviceCommandMappercommandMapper;AutowiredprivateRocketMQTemplaterocketMQTemplate;AutowiredprivateDeviceCommandServicecommandService;privatestaticfinalStringTIMEOUT_TOPICcommand-timeout-check;privatestaticfinalintMAX_RETRY2;/** * 下发指令时发送延迟消息延迟时间指令超时时间 */publicvoidsendTimeoutCheck(StringcommandId,intdelaySeconds){MessageStringmsgMessageBuilder.withPayload(commandId).build();// RocketMQ延迟级别: 1s1, 5s2, 10s3, 30s4, 1m5...intdelayLeveldelaySeconds5?2:(delaySeconds10?3:4);rocketMQTemplate.asyncSend(TIMEOUT_TOPIC,msg,newSendCallback(){OverridepublicvoidonSuccess(SendResultsendResult){}OverridepublicvoidonException(Throwablee){log.error(超时检查消息发送失败: commandId{},commandId,e);}},3000,delayLevel);}/** * 消费超时检查消息 */RocketMQMessageListener(topicTIMEOUT_TOPIC,consumerGroupcommand-timeout-group)ComponentpublicclassTimeoutConsumerimplementsRocketMQListenerString{OverridepublicvoidonMessage(StringcommandId){DeviceCommandcmdcommandMapper.selectByCommandId(commandId);if(cmdnull)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下发JSONObjectmessagebuildCommandMessage(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 服务端轮询接口RestControllerRequestMapping(/api/device)publicclassDevicePollController{AutowiredprivateDeviceCommandMappercommandMapper;/** * 设备拉取待执行指令 */GetMapping(/{sn}/commands/pending)publicResultListDeviceCommandgetPendingCommands(PathVariableStringsn){// 查询状态为已下发且未回调的指令ListDeviceCommandcommandscommandMapper.selectList(newLambdaQueryWrapperDeviceCommand().eq(DeviceCommand::getDeviceSn,sn).eq(DeviceCommand::getStatus,1).orderByAsc(DeviceCommand::getSendTime).last(LIMIT 5));returnResult.success(commands);}/** * 设备HTTP上报回调 */PostMapping(/{sn}/callback)publicResultVoidcallback(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){StringcommandIdcommand.getString(commandId);Stringcmdcommand.getString(command);longtimestampcommand.getLong(timestamp);Stringsigncommand.getString(sign);StringexpectedSignMD5Utils.md5(commandIdcmdtimestamptoken);returnexpectedSign.equals(sign);}9.3 防重放攻击设备端维护一个最近100条commandId的LRU缓存收到重复commandId直接忽略。配合timestamp字段超过5分钟的指令直接丢弃。十、通信协议完整定义表指令方向params回调data超时重试OPEN_DOOR下行orderId, maxDurationdoorOpen, openDuration10s2次CLOSE_DOOR下行无doorClosed10s1次RESTART下行delayrestartScheduled60s0次SYNC_TIME下行serverTimesynced, deviceTime5s1次INVENTORY下行无items[]30s1次UPDATE_CONFIG下行多个配置项updated10s1次UPLOAD_LOG下行startTime, endTimelogUrl60s0次TAKE_PHOTO下行cameraIdimageUrl10s1次十一、小结设备通信协议设计的核心就四个字可靠、幂等。commandId贯穿整个生命周期从下发到回调到超时检查全靠它串联。MQTT是主通道HTTP轮询是兜底RocketMQ延迟消息做超时检查三层保障确保指令不丢、不重、不卡。安全层面Token认证指令签名防重放三管齐下。这套协议跑通了设备端和服务端就能稳定对话后面的业务逻辑就是水到渠成的事。

相关新闻