Просмотр исходного кода

feat(devices): 统一设备离线检测机制并集成在线状态管理

- 在多个微服务中实现统一的设备离线检测定时任务
- 重构设备离线判断逻辑,支持可配置的超时时间和检查间隔
- 添加REST API方式更新设备在线状态到zk-api-service
- 实现设备状态变化时的API调用和日志记录功能
- 集成到各个消息处理器中自动标记设备在线状态
- 优化并发安全的数据结构管理设备状态信息
林仔 4 дней назад
Родитель
Сommit
baa3fd0edf
18 измененных файлов с 544 добавлено и 88 удалено
  1. 42 12
      audio-service/src/main/java/com/zksy/audio/utils/DeviceOfflineCheckTask.java
  2. 5 0
      audio-service/src/main/java/com/zksy/audio/utils/MessageHandler.java
  3. 79 0
      environment-service/src/main/java/com/zksy/environment/utils/DeviceOfflineCheckTask.java
  4. 55 12
      firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/DeviceOfflineCheckTask.java
  5. 7 3
      firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/MessageHandler.java
  6. 42 12
      flammable-gas-service/src/main/java/com/zksy/gas/utils/DeviceOfflineCheckTask.java
  7. 6 0
      flammable-gas-service/src/main/java/com/zksy/gas/utils/MessageHandler.java
  8. 54 11
      manhole-service/src/main/java/com/zksy/manhole/utils/DeviceOfflineCheckTask.java
  9. 5 1
      manhole-service/src/main/java/com/zksy/manhole/utils/MessageHandler.java
  10. 42 12
      radar-service/src/main/java/com/zksy/radar/utils/DeviceOfflineCheckTask.java
  11. 6 0
      radar-service/src/main/java/com/zksy/radar/utils/MessageHandler.java
  12. 42 12
      telemetry-service/src/main/java/com/zksy/telemetry/utils/DeviceOfflineCheckTask.java
  13. 7 0
      telemetry-service/src/main/java/com/zksy/telemetry/utils/MessageHandler.java
  14. 42 12
      water-level-service/src/main/java/com/zksy/water/utils/DeviceOfflineCheckTask.java
  15. 6 0
      water-level-service/src/main/java/com/zksy/water/utils/MessageHandler.java
  16. 35 0
      zk-api-service/src/main/java/com/zksy/api/controller/EquipmentStatusController.java
  17. 9 1
      zk-api-service/src/main/java/com/zksy/api/service/EquipmentStatusApiService.java
  18. 60 0
      zk-api-service/src/main/java/com/zksy/api/service/impl/EquipmentStatusApiServiceImpl.java

+ 42 - 12
audio-service/src/main/java/com/zksy/audio/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,66 @@
-package com.zksy.telemetry.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.audio.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 5 - 0
audio-service/src/main/java/com/zksy/audio/utils/MessageHandler.java

@@ -52,6 +52,8 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 	private final static String head = "524946462440010057415645666D742010000000010001000020000000400000020010006461746100400100";
 	private final NoiseInfoService service;
 	@Autowired
+	private DeviceOfflineCheckTask deviceOfflineCheckTask;
+	@Autowired
 	public MessageHandler(NoiseInfoService noiseInfoService) {
 		this.service = noiseInfoService;
 	}
@@ -288,6 +290,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 							DatePattern.NORM_DATETIME_PATTERN));
 					//log.info(JsonUtils.toJsonString(noise));
 					service.save(noise);
+					// 更新设备最后接收时间
+					DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(info.getDeviceNo(), new java.util.Date());
+					deviceOfflineCheckTask.markDeviceOnline(info.getDeviceNo());
 					//设备默认上报三次音频数据, 数据接收完成回复设备Ok,停止上报
 					byte[] cmdBytes = "OK\r\n".getBytes(StandardCharsets.US_ASCII);
 					ctx.channel().writeAndFlush(Unpooled.copiedBuffer(cmdBytes));

+ 79 - 0
environment-service/src/main/java/com/zksy/environment/utils/DeviceOfflineCheckTask.java

@@ -0,0 +1,79 @@
+package com.zksy.environment.utils;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
+
+import java.util.Date;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * 设备离线检测定时任务
+ * 通过 REST API 更新设备在线状态
+ */
+@Component
+public class DeviceOfflineCheckTask {
+    private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
+
+    /** 设备最后接收数据时间 */
+    public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    /** 已知离线设备集合,用于避免重复API调用 */
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
+    @Autowired
+    private RestTemplate restTemplate;
+
+    /** 离线判定超时,单位:分钟,默认30分钟 */
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
+
+    /** 检查间隔,单位:毫秒,默认5分钟 */
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
+    public void checkDeviceOffline() {
+        Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
+        for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
+            long diff = now.getTime() - entry.getValue().getTime();
+            if (diff > timeoutMs) {
+                // 只在设备不在离线集合中时才调用API
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
+            }
+        }
+    }
+
+    /**
+     * 标记设备在线(收到数据时调用)
+     * 只在设备当前处于离线状态时才调用API
+     */
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 55 - 12
firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,79 @@
-package com.zksy.pressure.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.pressure.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
+/**
+ * 设备离线检测定时任务
+ * 通过 REST API 更新设备在线状态
+ */
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
+    /** 设备最后接收数据时间 */
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    /** 已知离线设备集合,用于避免重复API调用 */
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    /** 离线判定超时,单位:分钟,默认30分钟 */
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    /** 检查间隔,单位:毫秒,默认5分钟 */
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                // 只在设备不在离线集合中时才调用API
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    /**
+     * 标记设备在线(收到数据时调用)
+     * 只在设备当前处于离线状态时才调用API
+     */
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 7 - 3
firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/MessageHandler.java

@@ -33,6 +33,7 @@ import org.springframework.web.client.RestTemplate;
 
 import java.math.BigDecimal;
 import java.time.LocalDateTime;
+import java.util.Date;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -56,6 +57,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
     @Autowired
     private RestTemplate restTemplate;
 
+    @Autowired
+    private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
     @Value("${device.status.api.url:http://zk-api-service/equipmentStatus/updateAlarmStatus}")
     private String deviceStatusApiUrl;
 	@Autowired
@@ -105,9 +109,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 				}
 
 				resultData.setId(UUID.randomUUID().toString());
-				// 更新设备最后一次接收数据的时间
-				//String addressCode = resultData.getAddressCode();
-				//DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(addressCode, new Date());
+				// 更新设备最后一次接收数据的时间,并标记设备在线
+				DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getTelemeteringStation(), new Date());
+				deviceOfflineCheckTask.markDeviceOnline(resultData.getTelemeteringStation());
 				firefightingPressureService.save(resultData);
 
 				// 4. 构建回复报文

+ 42 - 12
flammable-gas-service/src/main/java/com/zksy/gas/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,66 @@
-package com.zksy.gas.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.gas.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 6 - 0
flammable-gas-service/src/main/java/com/zksy/gas/utils/MessageHandler.java

@@ -58,6 +58,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
     @Autowired
     private RestTemplate restTemplate;
 
+    @Autowired
+    private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
     @Value("${device.status.api.url:http://zk-api-service/equipmentStatus/updateAlarmStatus}")
     private String deviceStatusApiUrl;
 
@@ -128,6 +131,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 			com.zksy.gas.domain.GasMonitorData resultData = DataParser.parseMessage(msgBytes);
 			resultData.setId(java.util.UUID.randomUUID().toString());
 			service.save(resultData);
+			// 更新设备最后接收时间,并标记设备在线
+			DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getMacAddress(), new java.util.Date());
+			deviceOfflineCheckTask.markDeviceOnline(resultData.getMacAddress());
 			checkIfSmsAlertNeeded(resultData);
 			logger.info("数据解析入库成功: {}", resultData);
 

+ 54 - 11
manhole-service/src/main/java/com/zksy/manhole/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,79 @@
 package com.zksy.manhole.utils;
 
-import com.zksy.manhole.service.BaseDevicesManholeService;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
+/**
+ * 设备离线检测定时任务
+ * 通过 REST API 更新设备在线状态
+ */
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
-    public static Map<String, Date> deviceLastReceiveTimeMap = new HashMap<>();
+    /** 设备最后接收数据时间 */
+    public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    /** 已知离线设备集合,用于避免重复API调用 */
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesManholeService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    /** 离线判定超时,单位:分钟,默认30分钟 */
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    /** 检查间隔,单位:毫秒,默认5分钟 */
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                // 只在设备不在离线集合中时才调用API
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
+
+    /**
+     * 标记设备在线(收到数据时调用)
+     * 只在设备当前处于离线状态时才调用API
+     */
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
 }

+ 5 - 1
manhole-service/src/main/java/com/zksy/manhole/utils/MessageHandler.java

@@ -24,6 +24,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 
 	private final ManholeDataService messageParseResultService;
 
+	@Autowired
+	private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
 	public MessageHandler(ManholeDataService messageParseResultService) {
 		this.messageParseResultService = messageParseResultService;
 	}
@@ -49,9 +52,10 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 				throw new InvalidMessageException("数据校验不成功");
 			} else {
 				ManholeData resultData = DataParser.parseMessage(msgString);
-				// 更新设备最后一次接收数据的时间
+				// 更新设备最后一次接收数据的时间,并标记设备在线
 				String imeiCardNumber = resultData.getImeiCardNumber();
 				DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(imeiCardNumber, new Date());
+				deviceOfflineCheckTask.markDeviceOnline(imeiCardNumber);
 				messageParseResultService.saveManholeData(resultData);
 			}
 		} catch (InvalidMessageException e) {

+ 42 - 12
radar-service/src/main/java/com/zksy/radar/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,66 @@
-package com.zksy.radar.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.radar.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 6 - 0
radar-service/src/main/java/com/zksy/radar/utils/MessageHandler.java

@@ -47,6 +47,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 	private final DevicePhoneFetchUtil devicePhoneFetchUtil;
     private final RestTemplate restTemplate;
 
+    @Autowired
+    private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
     @Value("${device.status.api.url:http://zk-api-service/equipmentStatus/updateAlarmStatus}")
     private String deviceStatusApiUrl;
 
@@ -108,6 +111,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 				}
 				// 入库
 				service.save(resultData);
+				// 更新设备最后接收时间,并标记设备在线
+				DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getSystemIdentifier(), new java.util.Date());
+				deviceOfflineCheckTask.markDeviceOnline(resultData.getSystemIdentifier());
 				logger.info("上报历史记录数据入库成功: {}", resultData);
 				
 				// 瞬时流量和流速告警入库

+ 42 - 12
telemetry-service/src/main/java/com/zksy/telemetry/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,66 @@
-package com.zksy.telemetry.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.telemetry.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 7 - 0
telemetry-service/src/main/java/com/zksy/telemetry/utils/MessageHandler.java

@@ -16,6 +16,8 @@ import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Component;
 
+import java.util.Date;
+
 @ChannelHandler.Sharable
 @Slf4j
 @Component
@@ -23,6 +25,8 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 	private static Logger logger = LoggerFactory.getLogger(MessageHandler.class);
 	private final TelemetryDataService service;
 	@Autowired
+	private DeviceOfflineCheckTask deviceOfflineCheckTask;
+	@Autowired
 	public MessageHandler(TelemetryDataService telemetryDataService) {
 		this.service = telemetryDataService;
 	}
@@ -71,6 +75,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 				}
 				// 入库
 				service.save(resultData);
+				// 更新设备最后接收时间,并标记设备在线
+				DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getSystemIdentifier(), new Date());
+				deviceOfflineCheckTask.markDeviceOnline(resultData.getSystemIdentifier());
 				logger.info("上报历史记录数据入库成功: {}", resultData);
 
 				// 如果结束位=1,可以根据需要在此处理关闭逻辑

+ 42 - 12
water-level-service/src/main/java/com/zksy/water/utils/DeviceOfflineCheckTask.java

@@ -1,36 +1,66 @@
-package com.zksy.radar.utils;/*
-package com.zksy.gasTransmitter.utils;
+package com.zksy.water.utils;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
 import org.springframework.scheduling.annotation.Scheduled;
 import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestTemplate;
 
 import java.util.Date;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
 @Component
 public class DeviceOfflineCheckTask {
     private static final Logger logger = LoggerFactory.getLogger(DeviceOfflineCheckTask.class);
 
-    // 用于存储设备编号及其最后一次接收数据的时间
     public static ConcurrentHashMap<String, Date> deviceLastReceiveTimeMap = new ConcurrentHashMap<>();
+
+    private final Set<String> offlineDeviceSet = ConcurrentHashMap.newKeySet();
+
     @Autowired
-    private BaseDevicesService baseDevicesService;
+    private RestTemplate restTemplate;
+
+    @Value("${device.offline.timeout-minutes:30}")
+    private int offlineTimeoutMinutes;
 
-    // 定时检查设备是否离线
-    @Scheduled(fixedRate = 24 * 60 * 60 * 1000) // 每24小时执行一次
+    @Scheduled(fixedRateString = "${device.offline.check-interval-ms:300000}")
     public void checkDeviceOffline() {
         Date now = new Date();
+        long timeoutMs = (long) offlineTimeoutMinutes * 60 * 1000;
         for (Map.Entry<String, Date> entry : deviceLastReceiveTimeMap.entrySet()) {
             long diff = now.getTime() - entry.getValue().getTime();
-            // 如果设备在 23 小时内没有接收数据,则认为设备离线
-            if (diff > 23 * 60 * 60 * 1000) {
-                baseDevicesService.getByDeviceNumberStatus(entry.getKey(), 1, 0);
-                logger.info("设备 {} 已离线", entry.getKey());
+            if (diff > timeoutMs) {
+                if (offlineDeviceSet.add(entry.getKey())) {
+                    updateDeviceOnlineStatus(entry.getKey(), 0);
+                    logger.info("设备 {} 已离线(超过{}分钟未收到数据)", entry.getKey(), offlineTimeoutMinutes);
+                }
             }
         }
     }
-}*/
+
+    public void markDeviceOnline(String deviceCode) {
+        if (offlineDeviceSet.remove(deviceCode)) {
+            updateDeviceOnlineStatus(deviceCode, 1);
+            logger.info("设备 {} 恢复在线", deviceCode);
+        }
+    }
+
+    private void updateDeviceOnlineStatus(String deviceCode, int onlineStatus) {
+        try {
+            Map<String, Object> params = Map.of("deviceCode", deviceCode, "onlineStatus", onlineStatus);
+            HttpHeaders headers = new HttpHeaders();
+            headers.setContentType(MediaType.APPLICATION_JSON);
+            HttpEntity<Map<String, Object>> request = new HttpEntity<>(params, headers);
+            restTemplate.postForObject("http://zk-api-service/equipmentStatus/updateOnlineStatus", request, Map.class);
+        } catch (Exception e) {
+            logger.error("更新设备在线状态失败: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+        }
+    }
+}

+ 6 - 0
water-level-service/src/main/java/com/zksy/water/utils/MessageHandler.java

@@ -47,6 +47,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 	private final DevicePhoneFetchUtil devicePhoneFetchUtil;
     private final RestTemplate restTemplate;
 
+    @Autowired
+    private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
     @Value("${device.status.api.url:http://zk-api-service/equipmentStatus/updateAlarmStatus}")
     private String deviceStatusApiUrl;
 
@@ -108,6 +111,9 @@ public class MessageHandler extends ChannelInboundHandlerAdapter {
 				}
 				// 入库
 				service.save(resultData);
+				// 更新设备最后接收时间,并标记设备在线
+				DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getSystemIdentifier(), new java.util.Date());
+				deviceOfflineCheckTask.markDeviceOnline(resultData.getSystemIdentifier());
 				logger.info("上报历史记录数据入库成功: {}", resultData);
 				
 				// 水位告警入库

+ 35 - 0
zk-api-service/src/main/java/com/zksy/api/controller/EquipmentStatusController.java

@@ -56,4 +56,39 @@ public class EquipmentStatusController {
             return Map.of("success", false, "msg", "更新失败: " + e.getMessage());
         }
     }
+
+    /**
+     * 更新设备在线状态
+     * @param params { deviceCode: "设备编码", onlineStatus: 0或1 }
+     * @return 更新结果
+     */
+    @PostMapping("/updateOnlineStatus")
+    @ApiOperation(value = "更新设备在线状态")
+    public Map<String, Object> updateOnlineStatus(@RequestBody Map<String, Object> params) {
+        String deviceCode = (String) params.get("deviceCode");
+        Integer onlineStatus = null;
+        Object onlineStatusObj = params.get("onlineStatus");
+        if (onlineStatusObj instanceof Integer) {
+            onlineStatus = (Integer) onlineStatusObj;
+        } else if (onlineStatusObj != null) {
+            onlineStatus = Integer.parseInt(onlineStatusObj.toString());
+        }
+
+        if (deviceCode == null || deviceCode.isEmpty()) {
+            log.warn("更新在线状态参数无效:deviceCode为空");
+            return Map.of("success", false, "msg", "deviceCode不能为空");
+        }
+        if (onlineStatus == null || (onlineStatus != 0 && onlineStatus != 1)) {
+            log.warn("更新在线状态参数无效:onlineStatus={}", onlineStatus);
+            return Map.of("success", false, "msg", "onlineStatus必须为0或1");
+        }
+
+        try {
+            boolean result = equipmentStatusApiService.updateOnlineStatus(deviceCode, onlineStatus);
+            return Map.of("success", result, "deviceCode", deviceCode, "onlineStatus", onlineStatus);
+        } catch (Exception e) {
+            log.error("更新设备在线状态异常: deviceCode={}, onlineStatus={}", deviceCode, onlineStatus, e);
+            return Map.of("success", false, "msg", "更新失败: " + e.getMessage());
+        }
+    }
 }

+ 9 - 1
zk-api-service/src/main/java/com/zksy/api/service/EquipmentStatusApiService.java

@@ -1,7 +1,7 @@
 package com.zksy.api.service;
 
 /**
- * 设备报警状态服务接口
+ * 设备状态服务接口
  */
 public interface EquipmentStatusApiService {
 
@@ -12,4 +12,12 @@ public interface EquipmentStatusApiService {
      * @return 是否更新成功
      */
     boolean updateAlarmStatus(String deviceCode, Integer alarmStatus);
+
+    /**
+     * 更新设备在线状态
+     * @param deviceCode 设备编码
+     * @param onlineStatus 在线状态:0-离线,1-在线
+     * @return 是否更新成功
+     */
+    boolean updateOnlineStatus(String deviceCode, Integer onlineStatus);
 }

+ 60 - 0
zk-api-service/src/main/java/com/zksy/api/service/impl/EquipmentStatusApiServiceImpl.java

@@ -95,4 +95,64 @@ public class EquipmentStatusApiServiceImpl implements EquipmentStatusApiService
 
         return true;
     }
+
+    @Override
+    @Transactional
+    public boolean updateOnlineStatus(String deviceCode, Integer onlineStatus) {
+        if (deviceCode == null || deviceCode.isEmpty()) {
+            log.warn("更新设备在线状态失败:deviceCode为空");
+            return false;
+        }
+
+        EquipmentBase equipmentBase = equipmentBaseApiMapper.selectOne(
+                new LambdaQueryWrapper<EquipmentBase>()
+                        .eq(EquipmentBase::getEquipmentCode, deviceCode)
+                        .last("LIMIT 1"));
+        if (equipmentBase == null) {
+            log.warn("更新设备在线状态失败:未找到设备编码对应的设备信息,deviceCode={}", deviceCode);
+            return false;
+        }
+
+        String equipmentId = equipmentBase.getEquipmentId();
+
+        EquipmentStatus existingStatus = equipmentStatusApiMapper.selectOne(
+                new LambdaQueryWrapper<EquipmentStatus>()
+                        .eq(EquipmentStatus::getEquipmentId, equipmentId));
+
+        if (existingStatus != null) {
+            Integer previousOnlineStatus = existingStatus.getOnlineStatus();
+            existingStatus.setOnlineStatus(onlineStatus);
+            existingStatus.setStatusUpdateTime(LocalDateTime.now());
+            int affected = equipmentStatusApiMapper.updateById(existingStatus);
+            if (affected <= 0) {
+                return false;
+            }
+            log.info("更新设备在线状态成功: deviceCode={}, equipmentId={}, onlineStatus={}",
+                    deviceCode, equipmentId, onlineStatus);
+            if (!Objects.equals(previousOnlineStatus, onlineStatus)) {
+                deviceStatusChangePublisher.publishAfterCommit(deviceCode, equipmentId,
+                        existingStatus.getCurrentStatus(), existingStatus.getAlarmStatus(),
+                        existingStatus.getOnlineStatus(), "zk-api:updateOnlineStatus");
+            }
+        } else {
+            EquipmentStatus newStatus = new EquipmentStatus();
+            newStatus.setEquipmentId(equipmentId);
+            newStatus.setOnlineStatus(onlineStatus);
+            newStatus.setAlarmStatus(0);
+            newStatus.setCurrentStatus(1);
+            newStatus.setStatusUpdateTime(LocalDateTime.now());
+            newStatus.setCreateTime(LocalDateTime.now());
+            int affected = equipmentStatusApiMapper.insert(newStatus);
+            if (affected <= 0) {
+                return false;
+            }
+            log.info("创建设备在线状态记录成功: deviceCode={}, equipmentId={}, onlineStatus={}",
+                    deviceCode, equipmentId, onlineStatus);
+            deviceStatusChangePublisher.publishAfterCommit(deviceCode, equipmentId,
+                    newStatus.getCurrentStatus(), newStatus.getAlarmStatus(),
+                    newStatus.getOnlineStatus(), "zk-api:updateOnlineStatus");
+        }
+
+        return true;
+    }
 }