fix (farm): 设备中心

- 添加redis缓存
This commit is contained in:
达富斌
2025-10-14 17:31:37 +08:00
parent af82b011ea
commit 9b918fadee
13 changed files with 686 additions and 26 deletions

View File

@@ -0,0 +1,30 @@
package cn.iocoder.yudao.module.farm.controller.admin.wx.vo;
import cn.iocoder.yudao.framework.common.pojo.PageParam;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* @description:
* @author: binda
* @date: 2025/9/19 15:35
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class DeviceInfoPageWxReqVO extends PageParam{
@Schema(description = "基地ID")
private String baseId;
@Schema(description = "设备类型ID")
private String deviceTypeId;
@Schema(description = "设备名称")
private String deviceName;
@Schema(description = "设备状态:0-未激活,1-在线,2-离线,3-维护中,4-故障")
private Integer deviceState;
@Schema(description = "绑定状态:0-未绑定; 1:已绑定")
private Integer bindingState;
}

View File

@@ -1,4 +1,4 @@
package cn.iocoder.yudao.module.farm.scheduled;
package cn.iocoder.yudao.module.farm.job;
import cn.iocoder.yudao.framework.quartz.core.handler.JobHandler;
import cn.iocoder.yudao.framework.tenant.core.job.TenantJob;

View File

@@ -1,4 +1,4 @@
package cn.iocoder.yudao.module.farm.scheduled;
package cn.iocoder.yudao.module.farm.job;
import cn.iocoder.yudao.framework.quartz.core.handler.JobHandler;
import cn.iocoder.yudao.framework.tenant.core.job.TenantJob;

View File

@@ -0,0 +1,27 @@
package cn.iocoder.yudao.module.farm.job;
import cn.iocoder.yudao.framework.quartz.core.handler.JobHandler;
import cn.iocoder.yudao.framework.tenant.core.job.TenantJob;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
/**
* @description: 土壤墒情设备预警
* @author: binda
* @date: 2025/10/14 14:21
*/
@Component
@TenantJob
@Slf4j
public class SoilDeviceWarningJob implements JobHandler {
@Override
public String execute(String param) throws Exception {
return "success";
}
}

View File

@@ -17,7 +17,7 @@ import java.util.List;
@NoArgsConstructor
public class RealTimeDataRespVO {
private String code;
private Integer code;
private String message;
private List<DeviceInfo> data;

View File

@@ -18,5 +18,6 @@ public class RealTimeDeviceDataReqVO {
* 设备编号 必填
*/
private String deviceId;
private Long tenantId;
}

View File

@@ -26,7 +26,7 @@ public class RealTimeDeviceDataRespVO {
@NoArgsConstructor
public static class DataRow {
private String id;
private String deviceId;
private String deviceId; // 设备唯一编号
private String property; // 重要参数
private String propertyName;
private String type;

View File

@@ -0,0 +1,186 @@
package cn.iocoder.yudao.module.farm.manager.equipmentproducer.redis;
import cn.iocoder.yudao.module.farm.enums.CommonConstants;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.soil.realtime.RealTimeDataRespVO;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.TimeUnit;
/**
* @description: 土壤墒情
* @author: binda
* @date: 2025/10/14 16:10
*/
@Slf4j
@Component
public class RedisFoeSoilRepository {
@Autowired
private StringRedisTemplate stringRedisTemplate;
@Autowired
private ObjectMapper objectMapper;
private static final String DEVICE_TYPE = CommonConstants.DEVICE_TYPE_SOIL;
/**
* 从Redis获取土壤设备实时数据
*/
public RealTimeDataRespVO getRealTimeData(Long tenantId, List<String> deviceAddrList) {
try {
if (deviceAddrList.isEmpty()) {
return null;
}
// 如果是单个设备,直接获取
if (deviceAddrList.size() == 1) {
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceAddrList.get(0));
String cachedJson = stringRedisTemplate.opsForValue().get(key);
if (cachedJson != null) {
return objectMapper.readValue(cachedJson, RealTimeDataRespVO.class);
}
return null;
}
// 多个设备,需要合并数据
List<String> keys = RedisKeyConstants.generateRealtimeDataKeys(DEVICE_TYPE, tenantId, deviceAddrList);
List<String> cachedJsons = stringRedisTemplate.opsForValue().multiGet(keys);
if (cachedJsons == null || cachedJsons.stream().allMatch(Objects::isNull)) {
return null;
}
// 合并所有设备数据
return mergeCachedData(cachedJsons, deviceAddrList);
} catch (Exception e) {
log.warn("从Redis获取土壤设备缓存数据失败: {}", e.getMessage());
return null;
}
}
/**
* 存储土壤设备实时数据到Redis
*/
public void storeRealTimeData(Long tenantId, RealTimeDataRespVO vendorResponse) {
if (vendorResponse == null || vendorResponse.getData() == null || vendorResponse.getData().isEmpty()) {
return;
}
try {
// 按设备分别存储
for (RealTimeDataRespVO.DeviceInfo deviceInfo : vendorResponse.getData()) {
String deviceAddr = deviceInfo.getDeviceAddr();
if (deviceAddr != null && !deviceAddr.trim().isEmpty()) {
// 创建单个设备的响应对象
RealTimeDataRespVO singleDeviceResponse = createSingleDeviceResponse(vendorResponse, deviceInfo);
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceAddr);
String value = objectMapper.writeValueAsString(singleDeviceResponse);
// 存储到Redis,设置6分钟过期时间
stringRedisTemplate.opsForValue().set(
key,
value,
RedisKeyConstants.DEVICE_REALTIME_DATA_TTL,
TimeUnit.SECONDS
);
log.debug("土壤设备实时数据已缓存,键: {}, 设备: {}", key, deviceAddr);
}
}
} catch (Exception e) {
log.warn("存储土壤设备实时数据到Redis失败: {}", e.getMessage());
}
}
/**
* 清理指定土壤设备的缓存
*/
public void clearDeviceCache(Long tenantId, String deviceAddr) {
try {
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceAddr);
Boolean deleted = stringRedisTemplate.delete(key);
if (Boolean.TRUE.equals(deleted)) {
log.info("已清理土壤设备缓存,键: {}", key);
}
} catch (Exception e) {
log.warn("清理土壤设备缓存失败: {}", e.getMessage());
}
}
/**
* 清理租户下所有土壤设备缓存
*/
public void clearTenantDeviceCache(Long tenantId) {
try {
String pattern = RedisKeyConstants.generateRealtimeDataPattern(DEVICE_TYPE, tenantId);
Set<String> keys = stringRedisTemplate.keys(pattern);
if (keys != null && !keys.isEmpty()) {
Long deletedCount = stringRedisTemplate.delete(keys);
log.info("已清理租户 {} 下所有土壤设备缓存,删除键数量: {}", tenantId, deletedCount);
}
} catch (Exception e) {
log.warn("清理租户土壤设备缓存失败,租户ID: {}, 错误: {}", tenantId, e.getMessage());
}
}
/**
* 合并多个设备的缓存数据
*/
private RealTimeDataRespVO mergeCachedData(List<String> cachedJsons, List<String> deviceAddrList) {
try {
RealTimeDataRespVO mergedResponse = new RealTimeDataRespVO();
mergedResponse.setCode(1000);
mergedResponse.setMessage("成功");
mergedResponse.setData(new ArrayList<>());
for (int i = 0; i < cachedJsons.size(); i++) {
String cachedJson = cachedJsons.get(i);
if (cachedJson != null) {
RealTimeDataRespVO singleResponse = objectMapper.readValue(cachedJson, RealTimeDataRespVO.class);
if (singleResponse != null && singleResponse.getData() != null) {
mergedResponse.getData().addAll(singleResponse.getData());
}
}
}
// 如果没有获取到任何有效数据,返回null
if (mergedResponse.getData().isEmpty()) {
return null;
}
return mergedResponse;
} catch (Exception e) {
log.warn("合并土壤设备缓存数据失败: {}", e.getMessage());
return null;
}
}
/**
* 创建单个设备的响应对象
*/
private RealTimeDataRespVO createSingleDeviceResponse(RealTimeDataRespVO originalResponse,
RealTimeDataRespVO.DeviceInfo deviceInfo) {
RealTimeDataRespVO singleResponse = new RealTimeDataRespVO();
singleResponse.setCode(originalResponse.getCode());
singleResponse.setMessage(originalResponse.getMessage());
List<RealTimeDataRespVO.DeviceInfo> singleDeviceList = new ArrayList<>();
singleDeviceList.add(deviceInfo);
singleResponse.setData(singleDeviceList);
return singleResponse;
}
}

View File

@@ -0,0 +1,217 @@
package cn.iocoder.yudao.module.farm.manager.equipmentproducer.redis;
import cn.iocoder.yudao.module.farm.enums.CommonConstants;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.wf.realtime.RealTimeDeviceDataRespVO;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* @description: 水肥机
* @author: binda
* @date: 2025/10/14 16:10
*/
@Slf4j
@Component
public class RedisFoeWfRepository {
@Autowired
private StringRedisTemplate stringRedisTemplate;
@Autowired
private ObjectMapper objectMapper;
private static final String DEVICE_TYPE = CommonConstants.DEVICE_TYPE_WF;
/**
* 从Redis获取水肥机设备实时数据
*/
public RealTimeDeviceDataRespVO getRealTimeData(Long tenantId, String deviceId) {
try {
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceId);
String cachedJson = stringRedisTemplate.opsForValue().get(key);
if (cachedJson != null) {
return objectMapper.readValue(cachedJson, RealTimeDeviceDataRespVO.class);
}
return null;
} catch (Exception e) {
log.warn("从Redis获取水肥机设备缓存数据失败: {}", e.getMessage());
return null;
}
}
/**
* 存储水肥机设备实时数据到Redis
*/
public void storeRealTimeData(Long tenantId, RealTimeDeviceDataRespVO vendorResponse) {
if (vendorResponse == null || vendorResponse.getData() == null || vendorResponse.getData().isEmpty()) {
return;
}
try {
// 从数据中提取设备ID(假设第一条数据的deviceId就是目标设备)
String deviceId = extractDeviceId(vendorResponse);
if (deviceId == null) {
log.warn("无法从水肥机实时数据中提取设备ID,跳过缓存");
return;
}
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceId);
String value = objectMapper.writeValueAsString(vendorResponse);
// 存储到Redis,设置6分钟过期时间
stringRedisTemplate.opsForValue().set(
key,
value,
RedisKeyConstants.DEVICE_REALTIME_DATA_TTL,
TimeUnit.SECONDS
);
log.debug("水肥机设备实时数据已缓存,键: {}, 设备ID: {}", key, deviceId);
} catch (Exception e) {
log.warn("存储水肥机设备实时数据到Redis失败: {}", e.getMessage());
}
}
/**
* 从响应数据中提取设备ID
*/
private String extractDeviceId(RealTimeDeviceDataRespVO vendorResponse) {
if (vendorResponse.getData() != null && !vendorResponse.getData().isEmpty()) {
// 返回第一条数据的设备ID
return vendorResponse.getData().get(0).getDeviceId();
}
return null;
}
/**
* 批量存储水肥机设备实时数据到Redis
* 用于处理多个设备的情况
*/
public void storeRealTimeDataBatch(Long tenantId, RealTimeDeviceDataRespVO vendorResponse) {
if (vendorResponse == null || vendorResponse.getData() == null || vendorResponse.getData().isEmpty()) {
return;
}
try {
// 按设备分组存储
Map<String, List<RealTimeDeviceDataRespVO.DataRow>> deviceDataMap = vendorResponse.getData().stream()
.filter(data -> data.getDeviceId() != null)
.collect(Collectors.groupingBy(RealTimeDeviceDataRespVO.DataRow::getDeviceId));
for (Map.Entry<String, List<RealTimeDeviceDataRespVO.DataRow>> entry : deviceDataMap.entrySet()) {
String deviceId = entry.getKey();
List<RealTimeDeviceDataRespVO.DataRow> deviceData = entry.getValue();
// 创建单个设备的响应对象
RealTimeDeviceDataRespVO singleDeviceResponse = createSingleDeviceResponse(vendorResponse, deviceData);
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceId);
String value = objectMapper.writeValueAsString(singleDeviceResponse);
// 存储到Redis
stringRedisTemplate.opsForValue().set(
key,
value,
RedisKeyConstants.DEVICE_REALTIME_DATA_TTL,
TimeUnit.SECONDS
);
log.debug("水肥机设备实时数据已缓存,键: {}, 设备ID: {}", key, deviceId);
}
} catch (Exception e) {
log.warn("批量存储水肥机设备实时数据到Redis失败: {}", e.getMessage());
}
}
/**
* 创建单个设备的响应对象
*/
private RealTimeDeviceDataRespVO createSingleDeviceResponse(RealTimeDeviceDataRespVO originalResponse,
List<RealTimeDeviceDataRespVO.DataRow> deviceData) {
RealTimeDeviceDataRespVO singleResponse = new RealTimeDeviceDataRespVO();
singleResponse.setCode(originalResponse.getCode());
singleResponse.setMsg(originalResponse.getMsg());
singleResponse.setData(deviceData);
return singleResponse;
}
/**
* 清理指定水肥机设备的缓存
*/
public void clearDeviceCache(Long tenantId, String deviceId) {
try {
String key = RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceId);
Boolean deleted = stringRedisTemplate.delete(key);
if (Boolean.TRUE.equals(deleted)) {
log.info("已清理水肥机设备缓存,键: {}", key);
}
} catch (Exception e) {
log.warn("清理水肥机设备缓存失败: {}", e.getMessage());
}
}
/**
* 清理租户下所有水肥机设备缓存
*/
public void clearTenantDeviceCache(Long tenantId) {
try {
String pattern = RedisKeyConstants.generateRealtimeDataPattern(DEVICE_TYPE, tenantId);
Set<String> keys = stringRedisTemplate.keys(pattern);
if (keys != null && !keys.isEmpty()) {
Long deletedCount = stringRedisTemplate.delete(keys);
log.info("已清理租户 {} 下所有水肥机设备缓存,删除键数量: {}", tenantId, deletedCount);
}
} catch (Exception e) {
log.warn("清理租户水肥机设备缓存失败,租户ID: {}, 错误: {}", tenantId, e.getMessage());
}
}
/**
* 获取多个水肥机设备的缓存数据
*/
public Map<String, RealTimeDeviceDataRespVO> getRealTimeDataBatch(Long tenantId, List<String> deviceIds) {
try {
if (deviceIds == null || deviceIds.isEmpty()) {
return Collections.emptyMap();
}
List<String> keys = deviceIds.stream()
.map(deviceId -> RedisKeyConstants.generateRealtimeDataKey(DEVICE_TYPE, tenantId, deviceId))
.collect(Collectors.toList());
List<String> cachedJsons = stringRedisTemplate.opsForValue().multiGet(keys);
Map<String, RealTimeDeviceDataRespVO> result = new HashMap<>();
for (int i = 0; i < deviceIds.size(); i++) {
String deviceId = deviceIds.get(i);
String cachedJson = cachedJsons != null ? cachedJsons.get(i) : null;
if (cachedJson != null) {
try {
RealTimeDeviceDataRespVO data = objectMapper.readValue(cachedJson, RealTimeDeviceDataRespVO.class);
result.put(deviceId, data);
} catch (Exception e) {
log.warn("解析水肥机设备 {} 缓存数据失败: {}", deviceId, e.getMessage());
}
}
}
return result;
} catch (Exception e) {
log.warn("批量获取水肥机设备缓存数据失败: {}", e.getMessage());
return Collections.emptyMap();
}
}
}

View File

@@ -0,0 +1,58 @@
package cn.iocoder.yudao.module.farm.manager.equipmentproducer.redis;
import java.util.List;
import java.util.stream.Collectors;
/**
* @description: Redis键
* @author: binda
* @date: 2025/9/30 11:22
*/
public class RedisKeyConstants {
/**
* 设备实时数据缓存前缀
* 格式: devices:deviceType:tenantId:deviceAddr
*/
public static final String DEVICE_REALTIME_DATA_KEY = "devices:%s:%s:%s";
/**
* 设备实时数据缓存过期时间(6分钟)
*/
public static final long DEVICE_REALTIME_DATA_TTL = 6 * 60;
/**
* 生成设备实时数据Redis键
* @param deviceType 设备类型
* @param tenantId 租户ID
* @param deviceAddr 设备地址
* @return Redis键
*/
public static String generateRealtimeDataKey(String deviceType, Long tenantId, String deviceAddr) {
return String.format(DEVICE_REALTIME_DATA_KEY, deviceType, tenantId, deviceAddr);
}
/**
* 生成设备实时数据Redis键(多个设备)
* @param deviceType 设备类型
* @param tenantId 租户ID
* @param deviceAddrs 设备地址列表
* @return Redis键列表
*/
public static List<String> generateRealtimeDataKeys(String deviceType, Long tenantId, List<String> deviceAddrs) {
return deviceAddrs.stream()
.map(deviceAddr -> generateRealtimeDataKey(deviceType, tenantId, deviceAddr))
.collect(Collectors.toList());
}
/**
* 生成设备实时数据Redis键模式(用于批量删除)
* @param deviceType 设备类型
* @param tenantId 租户ID
* @return Redis键模式
*/
public static String generateRealtimeDataPattern(String deviceType, Long tenantId) {
return String.format(DEVICE_REALTIME_DATA_KEY, deviceType, tenantId, "*");
}
}

View File

@@ -14,12 +14,12 @@ import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.soil.login.Use
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.soil.login.UserLoginRespVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.soil.realtime.RealTimeDataReqVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.soil.realtime.RealTimeDataRespVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.redis.RedisFoeSoilRepository;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.service.AbstractDeviceService;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.service.HistoryDataStorageService;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.service.TokenCacheService;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.utils.HttpClientUtils;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.util.*;
@@ -37,14 +37,13 @@ public class SoilService extends AbstractDeviceService {
private final HistoryDataStorageService historyDataStorageService;
private final DeviceConf deviceConf;
private final HttpClientUtils httpClientUtils;
private final StringRedisTemplate stringRedisTemplate;
private final RedisFoeSoilRepository redisForSoilRepository;
public SoilService(HistoryDataStorageService historyDataStorageService, DeviceConf deviceConf, HttpClientUtils httpClientUtils, StringRedisTemplate stringRedisTemplate) {
public SoilService(HistoryDataStorageService historyDataStorageService, DeviceConf deviceConf, HttpClientUtils httpClientUtils, RedisFoeSoilRepository redisForSoilRepository) {
this.historyDataStorageService = historyDataStorageService;
this.deviceConf = deviceConf;
this.httpClientUtils = httpClientUtils;
this.stringRedisTemplate = stringRedisTemplate;
this.redisForSoilRepository = redisForSoilRepository;
}
@Override
@@ -96,11 +95,30 @@ public class SoilService extends AbstractDeviceService {
// 转换请求参数
RealTimeDataReqVO soilRequest = (RealTimeDataReqVO) realTimeRequest;
Long tenantId = soilRequest.getTenantId();
String deviceAddrs = soilRequest.getDeviceAddrs();
// 解析设备地址列表
List<String> deviceAddrList = parseDeviceAddrs(deviceAddrs);
// 尝试从Redis获取缓存数据
RealTimeDataRespVO cachedData = redisForSoilRepository.getRealTimeData(tenantId, deviceAddrList);
if (cachedData != null) {
log.info("从Redis缓存获取土壤设备实时数据,租户ID: {}, 设备地址: {}", tenantId, deviceAddrs);
return RealTimeDataResponse.success(cachedData);
}
log.info("Redis缓存未命中,调用土壤设备厂商API获取实时数据,租户ID: {}, 设备地址: {}", tenantId, deviceAddrs);
// 调用土壤设备厂商API获取实时数据 GET 请求方式
RealTimeDataRespVO vendorResponse = callSoilRealTimeDataApi(token, soilRequest);
// 加入redis缓存
this.storeDataToRedis(soilRequest, vendorResponse);
if (vendorResponse != null && 1000 == vendorResponse.getCode()) {
// 保存到Redis缓存
redisForSoilRepository.storeRealTimeData(tenantId, vendorResponse);
log.info("土壤设备实时数据已缓存到Redis,租户ID: {}, 设备数量: {}",
tenantId, vendorResponse.getData() != null ? vendorResponse.getData().size() : 0);
}
// 直接返回原始响应数据,不进行统一转换
return RealTimeDataResponse.success(vendorResponse);
@@ -110,15 +128,6 @@ public class SoilService extends AbstractDeviceService {
}
}
private void storeDataToRedis(RealTimeDataReqVO soilRequest, RealTimeDataRespVO vendorResponse) {
if (vendorResponse == null || vendorResponse.getData() == null || vendorResponse.getData().isEmpty()) {
return;
}
Long tenantId = soilRequest.getTenantId();
}
@Override
protected HistoryDataResponse doGetHistoryData(String token, Object historyRequest) {
try {
@@ -251,7 +260,6 @@ public class SoilService extends AbstractDeviceService {
Map<String, String> headers = new HashMap<>();
headers.put("token", token);
// headers.put("Authorization", "Bearer " + token);
Map<String, Object> params = new HashMap<>();
if (request.getDeviceAddrs() != null) {
@@ -261,6 +269,34 @@ public class SoilService extends AbstractDeviceService {
return httpClientUtils.doGet(url, headers, params, RealTimeDataRespVO.class);
}
/**
* 解析设备地址字符串为列表
*/
private List<String> parseDeviceAddrs(String deviceAddrs) {
if (deviceAddrs == null || deviceAddrs.trim().isEmpty()) {
return Collections.emptyList();
}
return Arrays.stream(deviceAddrs.split(","))
.map(String::trim)
.filter(addr -> !addr.isEmpty())
.collect(Collectors.toList());
}
/**
* 清理土壤设备缓存(委托给Repository)
*/
public void clearDeviceCache(Long tenantId, String deviceAddr) {
redisForSoilRepository.clearDeviceCache(tenantId, deviceAddr);
}
/**
* 清理租户下所有土壤设备缓存(委托给Repository)
*/
public void clearTenantDeviceCache(Long tenantId) {
redisForSoilRepository.clearTenantDeviceCache(tenantId);
}
private HistoryDataListPageRespVO callSoilHistoryDataApi(String token, HistoryDataListPageReqVO request) {
String url = deviceConf.getFourSituations().getHosts() + deviceConf.getFourSituations().getSoil().getUrl44();

View File

@@ -14,16 +14,16 @@ import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.wf.realtime.Re
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.wf.realtime.RealTimeDeviceDataRespVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.wf.updateparams.UpdateDeviceStatusReqVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.dto.wf.updateparams.UpdateDeviceStatusRespVO;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.redis.RedisFoeWfRepository;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.service.AbstractDeviceService;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.service.TokenCacheService;
import cn.iocoder.yudao.module.farm.manager.equipmentproducer.utils.HttpClientUtils;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.text.MessageFormat;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.*;
import java.util.stream.Collectors;
/**
@@ -31,15 +31,18 @@ import java.util.stream.Collectors;
* @author: binda
* @date: 2025/9/25 20:42
*/
@Slf4j
@Service
public class WaterFertilizerService extends AbstractDeviceService {
private final DeviceConf deviceConf;
private final HttpClientUtils httpClientUtils;
private final RedisFoeWfRepository redisForWfRepository;
public WaterFertilizerService(DeviceConf deviceConf, HttpClientUtils httpClientUtils) {
public WaterFertilizerService(DeviceConf deviceConf, HttpClientUtils httpClientUtils, RedisFoeWfRepository redisForWfRepository) {
this.deviceConf = deviceConf;
this.httpClientUtils = httpClientUtils;
this.redisForWfRepository = redisForWfRepository;
}
@Override
@@ -93,10 +96,28 @@ public class WaterFertilizerService extends AbstractDeviceService {
try {
// 转换请求参数
RealTimeDeviceDataReqVO wfRequest = (RealTimeDeviceDataReqVO) realTimeRequest;
Long tenantId = wfRequest.getTenantId();
String deviceId = wfRequest.getDeviceId();
// 尝试从Redis获取缓存数据
RealTimeDeviceDataRespVO cachedData = redisForWfRepository.getRealTimeData(tenantId, deviceId);
if (cachedData != null) {
log.info("从Redis缓存获取水肥机设备实时数据,租户ID: {}, 设备ID: {}", tenantId, deviceId);
return RealTimeDataResponse.success(cachedData);
}
log.info("Redis缓存未命中,调用水肥机设备厂商API获取实时数据,租户ID: {}, 设备ID: {}", tenantId, deviceId);
// 调用水肥机设备厂商API获取实时数据 GET 请求方式
RealTimeDeviceDataRespVO vendorResponse = callWFRealTimeDataApi(token, wfRequest);
if (vendorResponse != null && "200".equals(vendorResponse.getCode())) {
// 保存到Redis缓存
redisForWfRepository.storeRealTimeData(tenantId, vendorResponse);
log.info("水肥机设备实时数据已缓存到Redis,租户ID: {}, 设备ID: {}", tenantId, deviceId);
}
// 直接返回原始响应数据,不进行统一转换
return RealTimeDataResponse.success(vendorResponse);
@@ -122,6 +143,58 @@ public class WaterFertilizerService extends AbstractDeviceService {
}
}
/**
* 批量获取水肥机设备实时数据
*/
public RealTimeDataResponse getRealTimeDataBatch(Long tenantId, String token, List<String> deviceIds) {
try {
// 尝试从Redis批量获取缓存数据
Map<String, RealTimeDeviceDataRespVO> cachedDataMap = redisForWfRepository.getRealTimeDataBatch(tenantId, deviceIds);
// 找出未命中的设备
List<String> missedDeviceIds = deviceIds.stream()
.filter(deviceId -> !cachedDataMap.containsKey(deviceId))
.collect(Collectors.toList());
// 如果所有设备都有缓存,直接返回
if (missedDeviceIds.isEmpty()) {
log.info("从Redis缓存批量获取水肥机设备实时数据,租户ID: {}, 设备数量: {}", tenantId, deviceIds.size());
// 合并所有设备数据
RealTimeDeviceDataRespVO mergedResponse = mergeDeviceResponses(cachedDataMap.values());
return RealTimeDataResponse.success(mergedResponse);
}
log.info("Redis缓存部分未命中,需要调用API的设备数量: {}", missedDeviceIds.size());
// 调用API获取未命中的设备数据
Map<String, RealTimeDeviceDataRespVO> apiDataMap = new HashMap<>();
for (String deviceId : missedDeviceIds) {
RealTimeDeviceDataReqVO request = new RealTimeDeviceDataReqVO();
request.setTenantId(tenantId);
request.setDeviceId(deviceId);
RealTimeDeviceDataRespVO apiResponse = callWFRealTimeDataApi(token, request);
if (apiResponse != null && "200".equals(apiResponse.getCode())) {
apiDataMap.put(deviceId, apiResponse);
// 缓存新获取的数据
redisForWfRepository.storeRealTimeData(tenantId, apiResponse);
}
}
// 合并缓存数据和API数据
Map<String, RealTimeDeviceDataRespVO> allDataMap = new HashMap<>(cachedDataMap);
allDataMap.putAll(apiDataMap);
RealTimeDeviceDataRespVO mergedResponse = mergeDeviceResponses(allDataMap.values());
return RealTimeDataResponse.success(mergedResponse);
} catch (Exception e) {
log.error("批量获取水肥机设备实时数据失败: {}", e.getMessage(), e);
return RealTimeDataResponse.error("批量获取水肥机设备实时数据失败: " + e.getMessage());
}
}
private LoginRespVO callWFVendorApi(String request) {
String url = deviceConf.getWaterAndFertilizer().getHosts() + deviceConf.getWaterAndFertilizer().getUrl1();
return httpClientUtils.doPost(url, null, request, LoginRespVO.class);
@@ -227,4 +300,36 @@ public class WaterFertilizerService extends AbstractDeviceService {
return UpdateDeviceResponse.error(vendorResponse.getMsg());
}
}
/**
* 清理水肥机设备缓存
*/
public void clearDeviceCache(Long tenantId, String deviceId) {
redisForWfRepository.clearDeviceCache(tenantId, deviceId);
}
/**
* 清理租户下所有水肥机设备缓存
*/
public void clearTenantDeviceCache(Long tenantId) {
redisForWfRepository.clearTenantDeviceCache(tenantId);
}
/**
* 合并多个设备的响应数据
*/
private RealTimeDeviceDataRespVO mergeDeviceResponses(Collection<RealTimeDeviceDataRespVO> responses) {
RealTimeDeviceDataRespVO mergedResponse = new RealTimeDeviceDataRespVO();
mergedResponse.setCode("200");
mergedResponse.setMsg("成功");
mergedResponse.setData(new ArrayList<>());
for (RealTimeDeviceDataRespVO response : responses) {
if (response != null && response.getData() != null) {
mergedResponse.getData().addAll(response.getData());
}
}
return mergedResponse;
}
}

View File

@@ -214,7 +214,7 @@ public class DeviceInfoServiceImpl implements DeviceInfoService {
// System.out.println("获取设备列表:"+response);
// 实时数据
RealTimeDeviceDataReqVO wfRequest = new RealTimeDeviceDataReqVO(deviceUniqueId);
RealTimeDeviceDataReqVO wfRequest = new RealTimeDeviceDataReqVO(deviceUniqueId,tenantId);
RealTimeDataResponse realTimeData = equipmentProducerManager.getRealTimeData(DEVICE_TYPE_WF, tenantId, token, wfRequest);
return new DeviceInfoDetailV2RespVO.WfInfo(deviceUniqueId,realTimeData);