【多线程】【实战】[示例]---- http接口堵塞,等待异步发送mqtt消息并拿到结果再接口响应
·
一、场景:
需求:安卓想查询固件最新版本。
流程:安卓调用http接口,后台接到请求后,发送mqtt消息给—>硬件设备,等设备返回返回最新版本信息给服务器端,服务器端再返回结果给安卓(超时时间为5秒)
二、接口服务层
@Getter
//用于等待MQTT返回结果,再返回HTTP接口响应 缓存:<设备SN , (等待锁 + 响应数据)>
public Map<String, DeviceResponseHolder> mqttResponseCache = new ConcurrentHashMap<>();
/**
* 【设备】获取设备固件版本版本号(通过发送mqtt消息给设备,异步获取固件版本号,堵塞5秒等待)
* @param deviceId 设备id
* @return ture 设备版本信息
*/
@Override
public String queryFirmwareVersionInfo(String snCode) {
//********************************变量************************************
//********************************校验************************************
//********************************业务************************************
//发送MQTT 获取固件当前版本,并更新数据库
//发送mqtt
String topic = "NZ/" + snCode + "/firmwareInfo/cmd";
//消息通用对象
MqttSendMessageDTO<DeviceFirmwareInfoMqttDTO> mqttMsg = new MqttSendMessageDTO<>();
//消息通用实体类内容略过.....
try {
log.debug("---------------------------------");
log.debug("【java】---发送获取固件信息 MQTT 消息");
mqttService.sendMqttMessage(topic,mqttMsg,2);
log.debug("---------------------------------");
}catch (Exception e) {
e.printStackTrace();
}
// 1. 创建等待锁(计数1)
CountDownLatch latch = new CountDownLatch(1);
DeviceResponseHolder holder = new DeviceResponseHolder(latch, null);
mqttResponseCache.put(snCode, holder);
boolean success = false;
String currentFirmwareVersion = null; //当前版本号
try {
// 2. 阻塞等待(最多等5秒,超时返回失败)
success = latch.await(5, TimeUnit.SECONDS); // 5秒内返回true; 5秒后返回false
AssertUtils.ifFalse(success, AppResponseCodeEnum.CODE_3010_DEVICE_ERROR); //异步查询固件版本号超时!
//成功返回
currentFirmwareVersion = holder.getCurrentFirmwareVersion();
devicePO.setFirmwareVersion(currentFirmwareVersion);
this.updateById(devicePO);//更新版本号
}catch (Exception e) {
e.printStackTrace();
}finally {
// 4. 清理缓存(防止内存泄漏)
mqttResponseCache.remove(snCode);
}
//获取最新固件版本
return currentFirmwareVersion;
}
三、线程处理类。
用于存放线程计数器,和异步返回的数据
package com.zgb.app.module.device.model.device;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.concurrent.CountDownLatch;
/**
* 响应持有类:封装等待锁和设备响应数据
*/
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true) // 自动生成无参构造器,强制初始化字段
public class DeviceResponseHolder {
//线程计数器
private final CountDownLatch latch;
//数据 (当前固件的版本号)异步获取
private String currentFirmwareVersion;
}
四、mqtt 处理类
接收到设备返回的消息,唤醒等待的线程
/**
* mqtt 服务器端接收到的全部消息
*/
private void handleMqttMessage(String topic, String payload) {
//【获取固件版本】■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■■
if(topic.contains("firmwareInfo/state")){
log.info("-----【获取固件版本】接收到mqtt获取固件版本的反馈--------");
// 使用 TypeReference 指定完整泛型类型
MqttSendMessageDTO<DeviceFirmwareInfoMqttDTO> dto = JSON.parseObject(payload, new TypeReference<>() {});
if(dto == null || dto.getData() == null){
log.error("-----【获取固件版本】接收到mqtt获取固件版本的反馈--------参数错误");
return ;
}
//获取消息里的字段
String snCode = dto.getData().getSnCode();
String firmwareVersion = dto.getData().getFirmwareVersion();
//获取阻塞map
Map<String, DeviceResponseHolder> mqttResponseCache = deviceServiceImpl.getMqttResponseCache();
DeviceResponseHolder deviceResponseHolder = mqttResponseCache.get(snCode);
if(deviceResponseHolder != null){
if(dto.getData()!=null){
deviceResponseHolder.setCurrentFirmwareVersion(firmwareVersion);
deviceResponseHolder.getLatch().countDown(); // 计数减1,唤醒阻塞
}
}
}
}
更多推荐
所有评论(0)