1
0

2 Коммиты de097d083d ... aeb0851fe6

Автор SHA1 Сообщение Дата
  林仔 aeb0851fe6 feat(devices): 统一设备离线检测与在线状态管理功能 4 дней назад
  林仔 baa3fd0edf feat(devices): 统一设备离线检测机制并集成在线状态管理 4 дней назад
20 измененных файлов с 554 добавлено и 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. 2 0
      environment-service/src/main/java/com/zksy/environment/EnvironmentApplication.java
  4. 8 0
      environment-service/src/main/java/com/zksy/environment/config/RSServerService.java
  5. 79 0
      environment-service/src/main/java/com/zksy/environment/utils/DeviceOfflineCheckTask.java
  6. 55 12
      firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/DeviceOfflineCheckTask.java
  7. 7 3
      firefighting-pressure-service/src/main/java/com/zksy/pressure/utils/MessageHandler.java
  8. 42 12
      flammable-gas-service/src/main/java/com/zksy/gas/utils/DeviceOfflineCheckTask.java
  9. 6 0
      flammable-gas-service/src/main/java/com/zksy/gas/utils/MessageHandler.java
  10. 54 11
      manhole-service/src/main/java/com/zksy/manhole/utils/DeviceOfflineCheckTask.java
  11. 5 1
      manhole-service/src/main/java/com/zksy/manhole/utils/MessageHandler.java
  12. 42 12
      radar-service/src/main/java/com/zksy/radar/utils/DeviceOfflineCheckTask.java
  13. 6 0
      radar-service/src/main/java/com/zksy/radar/utils/MessageHandler.java
  14. 42 12
      telemetry-service/src/main/java/com/zksy/telemetry/utils/DeviceOfflineCheckTask.java
  15. 7 0
      telemetry-service/src/main/java/com/zksy/telemetry/utils/MessageHandler.java
  16. 42 12
      water-level-service/src/main/java/com/zksy/water/utils/DeviceOfflineCheckTask.java
  17. 6 0
      water-level-service/src/main/java/com/zksy/water/utils/MessageHandler.java
  18. 35 0
      zk-api-service/src/main/java/com/zksy/api/controller/EquipmentStatusController.java
  19. 9 1
      zk-api-service/src/main/java/com/zksy/api/service/EquipmentStatusApiService.java
  20. 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));

+ 2 - 0
environment-service/src/main/java/com/zksy/environment/EnvironmentApplication.java

@@ -5,7 +5,9 @@ import org.mybatis.spring.annotation.MapperScan;
 import org.springframework.boot.SpringApplication;
 import org.springframework.boot.autoconfigure.SpringBootApplication;
 import org.springframework.scheduling.annotation.EnableAsync;
+import org.springframework.scheduling.annotation.EnableScheduling;
 
+@EnableScheduling
 @MapperScan({
         "com.zksy.environment.mapper",
         "com.zksy.base.mapper",

+ 8 - 0
environment-service/src/main/java/com/zksy/environment/config/RSServerService.java

@@ -6,6 +6,7 @@ import com.zksy.api.utils.SmsUtil;
 import com.zksy.environment.domain.ERealTimeData;
 import com.zksy.environment.mapper.ERealTimeDataMapper;
 import com.zksy.environment.utils.AlarmUtil;
+import com.zksy.environment.utils.DeviceOfflineCheckTask;
 import com.zksy.utils.DevicePhoneFetchUtil;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.beans.factory.annotation.Autowired;
@@ -22,6 +23,7 @@ import javax.annotation.PreDestroy;
 import java.math.BigDecimal;
 import java.text.SimpleDateFormat;
 import java.time.LocalDateTime;
+import java.util.Date;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -58,6 +60,9 @@ public class RSServerService {
     @Autowired
     private RestTemplate restTemplate;
 
+    @Autowired
+    private DeviceOfflineCheckTask deviceOfflineCheckTask;
+
     @Value("${device.status.api.url:http://zk-api-service/equipmentStatus/updateAlarmStatus}")
     private String deviceStatusApiUrl;
 
@@ -130,6 +135,9 @@ public class RSServerService {
                 @Override
                 public void receiveRealtimeData(RealTimeData data) {
                     String deviceId = String.valueOf(data.getDeviceId());
+                    // 更新设备最后接收数据时间,并标记设备在线
+                    DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(deviceId, new Date());
+                    deviceOfflineCheckTask.markDeviceOnline(deviceId);
                     boolean hasAlarm = false;
                     for (NodeData nd : data.getNodeList()) {
                         try {

+ 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;
+    }
 }