From f453f179039f0599e09921da5f7a049baf5bf6e7 Mon Sep 17 00:00:00 2001 From: gaoyutao Date: Fri, 27 Feb 2026 17:49:01 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96ip=E5=85=AC=E7=BD=91=E6=8E=A2?= =?UTF-8?q?=E6=B5=8B=E9=99=8D=E7=BA=A7=E9=80=BB=E8=BE=91=E3=80=82=20?= =?UTF-8?q?=E4=BC=98=E5=8C=96=E7=BD=91=E7=BB=9C=E6=B5=81=E9=87=8F=E9=87=8D?= =?UTF-8?q?=E8=AF=95=E4=B8=AA=E5=88=AB=E6=96=87=E4=BB=B6=E5=8F=91=E9=80=81?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5=E9=80=BB=E8=BE=91=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduler/config/SchedulerConfig.java | 17 ++++ .../scheduler/service/BusinessTasks.java | 2 +- .../client/service/impl/AgentServiceImpl.java | 94 +++++++++++++------ .../agent/client/utils/PublicIpFetcher.java | 3 +- 4 files changed, 83 insertions(+), 33 deletions(-) diff --git a/src/main/java/com/tongran/agent/client/scheduler/config/SchedulerConfig.java b/src/main/java/com/tongran/agent/client/scheduler/config/SchedulerConfig.java index 928dde6..f1879da 100644 --- a/src/main/java/com/tongran/agent/client/scheduler/config/SchedulerConfig.java +++ b/src/main/java/com/tongran/agent/client/scheduler/config/SchedulerConfig.java @@ -65,4 +65,21 @@ public class SchedulerConfig { executor.initialize(); return executor; } + + // 新增:心跳专用线程池 + @Bean("heartbeatExecutor") + public ThreadPoolTaskExecutor heartbeatExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(2); + executor.setMaxPoolSize(4); + executor.setQueueCapacity(0); // 重要:设为0,不要队列! + executor.setKeepAliveSeconds(30); + executor.setThreadNamePrefix("heartbeat-"); + + // 重要:使用CallerRunsPolicy,确保心跳一定执行 + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + + executor.initialize(); + return executor; + } } \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/client/scheduler/service/BusinessTasks.java b/src/main/java/com/tongran/agent/client/scheduler/service/BusinessTasks.java index 9f373fb..7312400 100644 --- a/src/main/java/com/tongran/agent/client/scheduler/service/BusinessTasks.java +++ b/src/main/java/com/tongran/agent/client/scheduler/service/BusinessTasks.java @@ -119,7 +119,7 @@ public class BusinessTasks { /** * 任务1:心跳上报任务 */ - @Async("taskExecutor") + @Async("heartbeatExecutor") public void heartbeatTask() { long timestamp = System.currentTimeMillis(); timestamp = Math.round(timestamp / 1000.0); diff --git a/src/main/java/com/tongran/agent/client/service/impl/AgentServiceImpl.java b/src/main/java/com/tongran/agent/client/service/impl/AgentServiceImpl.java index def0038..fe2e4d5 100644 --- a/src/main/java/com/tongran/agent/client/service/impl/AgentServiceImpl.java +++ b/src/main/java/com/tongran/agent/client/service/impl/AgentServiceImpl.java @@ -2454,44 +2454,76 @@ public class AgentServiceImpl implements AgentService { } private void processSingleFile(File file, boolean lastTraffic) { - // 使用Properties类读取,并指定UTF-8编码 - Properties props = new Properties(); - try (InputStream input = Files.newInputStream(file.toPath()); - InputStreamReader reader = new InputStreamReader(input, StandardCharsets.UTF_8)) { // 关键:指定UTF-8编码 - - props.load(reader); // 使用指定编码的Reader加载 - String data = props.getProperty("traffic"); - if (StringUtils.isNotBlank(data)) { - List list = JSONArray.parseArray(data, NetVO.class); - if(list != null){ - for (NetVO netVO : list) { - netVO.setLastTrafficFlag(lastTraffic); + // 重试配置 + int maxRetry = 3; // 最大重试次数 + for (int retryCount = 0; retryCount < maxRetry; retryCount++) { + try { + // 每次重试前都检查连接 + boolean success = true; + int activeConnect = client.getActiveConnections(); + if (activeConnect == 0) { + success = connection(); + if (!success) { + AssertLog.info("连接失败,第 {} 次重试,文件:{}", + retryCount + 1, file.getName()); + continue; } } - data = JSONObject.toJSONString(list); - Message message = Message.builder() - .clientId(GlobalConfig.CLIENT_ID) - .dataType(MsgEnum.网络上报重试.getValue()) - .data(data) - .build(); + // 使用Properties类读取,并指定UTF-8编码 + Properties props = new Properties(); + try (InputStream input = Files.newInputStream(file.toPath()); + InputStreamReader reader = new InputStreamReader(input, StandardCharsets.UTF_8)) { // 关键:指定UTF-8编码 - sessionManager.writeAndFlush( - sessionManager.getSessionById(GlobalConfig.CLIENT_ID).getChannel(), - message - ); + props.load(reader); // 使用指定编码的Reader加载 + String data = props.getProperty("traffic"); + if (StringUtils.isNotBlank(data)) { + List list = JSONArray.parseArray(data, NetVO.class); + if(list != null){ + for (NetVO netVO : list) { + netVO.setLastTrafficFlag(lastTraffic); + } + } + data = JSONObject.toJSONString(list); + Message message = Message.builder() + .clientId(GlobalConfig.CLIENT_ID) + .dataType(MsgEnum.网络上报重试.getValue()) + .data(data) + .build(); - if (lastTraffic) { - AssertLog.info("发送最后一条重试网络信息包={}", JSON.toJSONString(message)); - } else { - AssertLog.info("发送重试recover网络信息包={}", JSON.toJSONString(message)); + sessionManager.writeAndFlush( + sessionManager.getSessionById(GlobalConfig.CLIENT_ID).getChannel(), + message + ); + + if (lastTraffic) { + AssertLog.info("发送最后一条重试网络信息包={}", JSON.toJSONString(message)); + } else { + AssertLog.info("发送重试recover网络信息包={}", JSON.toJSONString(message)); + } + return; // 成功返回 + } else { + AssertLog.info("文件中没有有效数据,删除文件:{}", file.getName()); + file.delete(); + return; + } + } + }catch (Exception e) { + AssertLog.error("loadNetFile异常,文件:{}, 重试次数:{}, 异常信息:{}", + file.getName(), retryCount + 1, e.getMessage()); + + // 如果不是最后一次重试,等待2秒再试 + if (retryCount < maxRetry - 1) { + try { + Thread.sleep(2000); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + return; + } } - } else { - AssertLog.info("文件中没有有效数据,删除文件:{}", file.getName()); - file.delete(); } - } catch (Exception e) { - AssertLog.error("loadNetFile异常,文件:{},异常信息:{}", file.getName(), e.getMessage()); } + // 所有重试都失败 + AssertLog.error("文件处理失败,已达最大重试次数({}),文件:{}", maxRetry, file.getName()); } public void caseTypeBySystem(String type, int interval, boolean collect){ diff --git a/src/main/java/com/tongran/agent/client/utils/PublicIpFetcher.java b/src/main/java/com/tongran/agent/client/utils/PublicIpFetcher.java index b57cb46..fe2a025 100644 --- a/src/main/java/com/tongran/agent/client/utils/PublicIpFetcher.java +++ b/src/main/java/com/tongran/agent/client/utils/PublicIpFetcher.java @@ -215,7 +215,7 @@ public class PublicIpFetcher { return result; } catch (Exception e) { AssertLog.error("HTTP主服务请求失败: {}", e.getMessage()); - return tryHttpFallback(); + throw e; // 不在这里降级,只抛异常 } }); @@ -224,6 +224,7 @@ public class PublicIpFetcher { } catch (TimeoutException e) { AssertLog.error("HTTP请求超时"); future.cancel(true); + return tryHttpFallback(); // 主线程执行降级 } catch (Exception e) { AssertLog.error("HTTP请求异常: {}", e.getMessage()); } finally {