#优化更新 连接上位机相关

This commit is contained in:
2026-09-09 09:01:54 +08:00
parent fda1060b38
commit 4ee3f6be90
21 changed files with 330 additions and 218 deletions

View File

@@ -1,29 +1,19 @@
package com.maibu.service;
import cn.hutool.core.io.IORuntimeException;
import cn.hutool.core.io.IoUtil;
import com.alibaba.fastjson2.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.maibu.common.NettyCacheKey;
import com.maibu.constant.Constant;
import com.maibu.core.business.Device;
import com.maibu.core.business.DeviceRunningStatusHistory;
import com.maibu.core.business.DeviceStatusRecordDTO;
import com.maibu.core.business.ErrorIdentificationStandard;
import com.maibu.core.business.alarm_center.AlarmMessage;
import com.maibu.core.business.device.NettyDevice;
import com.maibu.core.business.inter.WebsocketMesDispather;
import com.maibu.core.enums.*;
import com.maibu.core.host.HeartBeatDTO;
import com.maibu.core.redis.RedisCache;
import com.maibu.dto.DeviceErrorPushDTO;
import com.maibu.mapper.ErrorIdentificationStandardMapper;
import com.maibu.memory.DeviceSessionManager;
import com.maibu.memory.GlobalMemory;
import com.maibu.memory.SiteMemory;
import com.maibu.mqtt.MqttTopic;
import com.maibu.utils.CompareUtils;
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.lang.reflect.Field;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
@@ -32,17 +22,37 @@ import org.springframework.core.io.ResourceLoader;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import javax.annotation.PreDestroy;
import java.io.IOException;
import java.lang.reflect.Field;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import com.alibaba.fastjson2.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.maibu.constant.Constant;
import com.maibu.constant.NettyCacheKey;
import com.maibu.core.business.Device;
import com.maibu.core.business.DeviceRunningStatusHistory;
import com.maibu.core.business.DeviceStatusRecordDTO;
import com.maibu.core.business.ErrorIdentificationStandard;
import com.maibu.core.business.alarm_center.AlarmMessage;
import com.maibu.core.business.device.NettyDevice;
import com.maibu.core.business.inter.WebsocketMesDispather;
import com.maibu.core.enums.CompareEnum;
import com.maibu.core.enums.ConnectorStatus;
import com.maibu.core.enums.DeviceStatus;
import com.maibu.core.enums.ErrorSource;
import com.maibu.core.enums.RespondCode;
import com.maibu.core.host.HeartBeatDTO;
import com.maibu.core.host.LocationMessage;
import com.maibu.core.redis.RedisCache;
import com.maibu.dto.DeviceErrorPushDTO;
import com.maibu.mapper.ErrorIdentificationStandardMapper;
import com.maibu.mapper.HostLocationMapper;
import com.maibu.memory.DeviceSessionManager;
import com.maibu.memory.GlobalMemory;
import com.maibu.memory.SiteMemory;
import com.maibu.mqtt.MqttTopic;
import com.maibu.utils.CompareUtils;
import cn.hutool.core.io.IORuntimeException;
import cn.hutool.core.io.IoUtil;
import lombok.extern.slf4j.Slf4j;
@Slf4j
@Service
@@ -75,10 +85,13 @@ public class DeviceThreadService {
@Autowired
private ScheduledExecutorService executor;
@Autowired
private HostLocationMapper hostLocationMapper;
public void deviceErrorMonitor() throws InterruptedException, IOException {
initStandard();
// TODO fix 后续是否分多个线程 一个site一个线程
// TODO fix 后续是否分多个线程 一个site一个线程
// 有还需要更新接受故障方式。
executor.scheduleWithFixedDelay(() -> {
try {
@@ -188,8 +201,10 @@ public class DeviceThreadService {
// }
// }
try {
GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId, String.format(MqttTopic.DEVICE_ERROR_PUSH_TOPIC, deviceId), dto);
// GlobalMemory.mqttClientUtil.publish(String.format(MqttTopic.DEVICE_ERROR_PUSH_TOPIC, deviceId), dto);
GlobalMemory.customMqttDeviceMonitor.publish(mqttClientId,
String.format(MqttTopic.DEVICE_ERROR_PUSH_TOPIC, deviceId), dto);
// GlobalMemory.mqttClientUtil.publish(String.format(MqttTopic.DEVICE_ERROR_PUSH_TOPIC,
// deviceId), dto);
List<NettyDevice> controlMasters = sessionManager.getAllSlaveControl(deviceId);
if (!CollectionUtils.isEmpty(controlMasters)) {
controlMasters.forEach(x -> {
@@ -203,7 +218,7 @@ public class DeviceThreadService {
}
public List<AlarmMessage> compareErrorStandard(Long siteId, Long orgId, DeviceRunningStatusHistory history,
List<ErrorIdentificationStandard> standards) {
List<ErrorIdentificationStandard> standards) {
if (history == null || CollectionUtils.isEmpty(standards))
return null;
List<AlarmMessage> list = new ArrayList<>();
@@ -212,7 +227,8 @@ public class DeviceThreadService {
String compareValues = standard.getCompareValues();
CompareEnum compareType = standard.getCompareType();
String targetValue = knownClassFieldFinding(fieldName, history);
log.debug("compare2 compareType:{} compareValues:{},targetValue:{}", compareType, compareValues, targetValue);
log.debug("compare2 compareType:{} compareValues:{},targetValue:{}", compareType, compareValues,
targetValue);
if (!StringUtils.isEmpty(compareValues) && !StringUtils.isEmpty(targetValue)) {
boolean result = CompareUtils.compare(compareType, targetValue, compareValues);
log.debug("compare3 result:{}", result);
@@ -253,7 +269,6 @@ public class DeviceThreadService {
return null;
}
/**
* 监控心跳
*/
@@ -268,10 +283,22 @@ public class DeviceThreadService {
Long lastActiveTime = x.getSysTimestamp();
NettyDevice nettyDevice = sessionManager.getDevice(x.getSn());
if (nowTime - lastActiveTime >= Constant.hostOfflineInterval) {
//超时了判断离线
// 超时了判断离线
if (nettyDevice != null) {
nettyDevice.status(0);
//todo 推送设备离线
nettyDevice.status(DeviceStatus.offline.getType());
// todo 推送设备离线
String key = NettyCacheKey.latestPositonKey + deviceId;
// 保存一下最新一次的数据到 数据库
LocationMessage locationMessage = redisCache.getCacheObject(key);
if (locationMessage != null) {
LambdaQueryWrapper<LocationMessage> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(LocationMessage::getDeviceId, deviceId);
LocationMessage existing = hostLocationMapper.selectOne(queryWrapper);
if (existing != null) {
locationMessage.setId(existing.getId());
}
hostLocationMapper.saveOrUpdate(locationMessage);
}
}
} else {
if (nettyDevice == null) {
@@ -289,31 +316,36 @@ public class DeviceThreadService {
device.setProductName("割草机产品MC700");
device.setCreateTime(LocalDateTime.now());
device.setUpdateTime(LocalDateTime.now());
device.setHasHost(true);
deviceService.insertDeviceBySelf(device);
} else {
device.setStatus(3);
device.setOnlineStatus(DeviceStatus.online.getType());
device.setUpdateTime(LocalDateTime.now());
if (!device.isHasHost()) {
device.setHasHost(true);
}
deviceService.updateDeviceBySelf(device);
}
// 推送json 找到设备绑定的主机进行推送
// List<NettyDevice> controlMasters = sessionManager.getAllSlaveControl(device.getSerialNumber());
// if (!CollectionUtils.isEmpty(controlMasters)) {
// controlMasters.forEach(x -> {
// DeviceStatusChangeDTO changeDTO = new DeviceStatusChangeDTO();
// changeDTO.setDeviceId(data);
// changeDTO.setStatus(DeviceStatus.online.getCode());
// changeDTO.setEvent(RespondCode.device_status_changed);
// ByteBuf buf = ctx.alloc().buffer();
// CommandUtils.buildCommand(buf, changeDTO, CommandConstant.interaction);
// String content = buf.toString(CharsetUtil.UTF_8);
// logger.info("推送下位机登录状态消息:{}", content);
// if (x.getChannel() != null && x.getChannel().isActive()) {
// x.getChannel().writeAndFlush(buf);
// }
// });
// List<NettyDevice> controlMasters =
// sessionManager.getAllSlaveControl(device.getSerialNumber());
// if (!CollectionUtils.isEmpty(controlMasters)) {
// controlMasters.forEach(x -> {
// DeviceStatusChangeDTO changeDTO = new DeviceStatusChangeDTO();
// changeDTO.setDeviceId(data);
// changeDTO.setStatus(DeviceStatus.online.getCode());
// changeDTO.setEvent(RespondCode.device_status_changed);
// ByteBuf buf = ctx.alloc().buffer();
// CommandUtils.buildCommand(buf, changeDTO, CommandConstant.interaction);
// String content = buf.toString(CharsetUtil.UTF_8);
// logger.info("推送下位机登录状态消息:{}", content);
// if (x.getChannel() != null && x.getChannel().isActive()) {
// x.getChannel().writeAndFlush(buf);
// }
// });
} else {
//todo 推送设备上线
// todo 推送设备上线
nettyDevice.status(1);
}
}