重新部署一下

This commit is contained in:
2023-03-23 09:14:28 +08:00
parent 04fb8d3388
commit f2d313721f
2 changed files with 153 additions and 32 deletions
@@ -61,7 +61,7 @@ public class CheckInstancesStatus {
int corePoolSize = Math.min(thread + 1, containerCount); int corePoolSize = Math.min(thread + 1, containerCount);
int maxPoolSize = Math.max(thread + 1, containerCount); int maxPoolSize = Math.max(thread + 1, containerCount);
// 创建一个包含10个线程的线程池 // 创建一个线程池
ExecutorService executor = new ThreadPoolExecutor( ExecutorService executor = new ThreadPoolExecutor(
corePoolSize, corePoolSize,
maxPoolSize, maxPoolSize,
@@ -412,16 +412,102 @@ class NoodlesApplicationTests {
public void checkVmStatusMultiThread() throws InterruptedException { public void checkVmStatusMultiThread() throws InterruptedException {
long startTime = System.nanoTime(); // 记录开始时间 long startTime = System.nanoTime(); // 记录开始时间
ExecutorService executorService = Executors.newFixedThreadPool(10); // 创建一个线程池 List<Container> containers = containerService.list();
List<HostMachine> vms = hostMachineService.list(); // 获取所有主机列表 int thread = CpuUtil.getLogicProcessorCount();
CountDownLatch countDownLatch = new CountDownLatch(vms.size()); // 用于等待所有线程完成 int containerCount = containerService.count();
for (HostMachine vm : vms) { int corePoolSize = Math.min(thread + 1, containerCount);
executorService.submit(() -> { int maxPoolSize = Math.max(thread + 1, containerCount);
ExecutorService executor = new ThreadPoolExecutor(
corePoolSize,
maxPoolSize,
1,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(),
new ThreadPoolExecutor.AbortPolicy());
// 创建一个 Future 列表,用于存储每个容器检查的结果
List<Future<Container>> futures = new ArrayList<>();
for (Container container : containers) {
futures.add(executor.submit(() -> {
try { try {
if (!StringUtils.hasText(vm.getManageIp())) { if (!StringUtils.hasText(container.getWebUi()) && !StringUtils.hasText(container.getServerAddress())) {
container.setContainerState(SystemConstant.CONTAINER_STATE_UNKNOWN);
return container;
}
//通过Socket(IP+端口)判断是否在线
if (StringUtils.hasText(container.getServerAddress())) {
if (UrlUtil.isServiceOnline(UrlUtil.getHostname(container.getServerAddress()), UrlUtil.getPort(container.getServerAddress()))) {
container.setContainerState(SystemConstant.CONTAINER_STATE_RUNNING);
return container;
}
}
//通过HTTP请求判断是否在线
HttpRequest.get(container.getWebUi()).setConnectionTimeout(5000).execute(true);
log.info("[实例状态检测]:与{}建立连接成功", container.getName());
container.setContainerState(SystemConstant.CONTAINER_STATE_RUNNING);
} catch (Exception exception) {
log.info("[实例状态检测]:与{}建立连接失败,状态转为离线", container.getName());
container.setContainerState(SystemConstant.CONTAINER_STATE_EXITED);
// 判断 Redis 里面有没有,如果有就不需要提醒,如果没有就提醒
Set<String> set = redisCache.getCacheSet(RedisConstant.NOTIFY_CONTAINER_IDS);
if (!CollectionUtils.isEmpty(set)) {
// 如果 Redis 里面有,说明已经发送过了未 check 的通知,所以不需要发送,直接返回
if (set.contains(container.getId().toString())) {
return container;
} else {
// 根据实例对应的提醒方式进行提醒
if (container.getNotify().equals(SystemConstant.NOTIFY_NO)) {
return container;
} else {
notificationService.sendContainerOfflineNotification(container.getId());
}
}
}
// 存入 Redis
set.add(container.getId().toString());
redisCache.setCacheSet(RedisConstant.NOTIFY_CONTAINER_IDS, set);
}
return container;
}));
}
// 等待所有线程执行完成,并收集更新后的容器列表
List<Container> updatedContainers = new ArrayList<>();
for (Future<Container> future : futures) {
try {
Container container = future.get();
updatedContainers.add(container);
} catch (InterruptedException | ExecutionException e) {
log.error("检查容器状态时出错:{}", e.getMessage());
}
}
// 将更新后的容器列表保存到数据库中
containerService.saveOrUpdateBatch(updatedContainers);
// 关闭线程池
executor.shutdown();
//检测虚拟机在线状态
LambdaQueryWrapper<HostMachine> vmWrapper = new LambdaQueryWrapper<>();
vmWrapper.ne(HostMachine::getHostMachineId, SystemConstant.HOST_MACHINE_ID_HOST);
List<HostMachine> vms = hostMachineService.list(vmWrapper);
vms.forEach(vm -> {
try {
if (!StringUtils.hasText(vm.getServerAddress()) && !StringUtils.hasText(vm.getServerAddress())) {
vm.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_UNKNOWN); vm.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_UNKNOWN);
return; return;
} }
//通过Ping的方式判断是否在线
if (StringUtils.hasText(vm.getServerAddress())) {
if (UrlUtil.isHostOnline(vm.getServerAddress())){
log.info("[实例状态检测]:与{}建立连接成功", vm.getName());
vm.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_ONLINE);
return;
}
}
//通过HTTP请求的方式判断是否在线
HttpRequest.get(vm.getManageIp()).setConnectionTimeout(5000).execute(true); HttpRequest.get(vm.getManageIp()).setConnectionTimeout(5000).execute(true);
log.info("[实例状态检测]:与{}建立连接成功", vm.getName()); log.info("[实例状态检测]:与{}建立连接成功", vm.getName());
vm.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_ONLINE); vm.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_ONLINE);
@@ -446,14 +532,49 @@ class NoodlesApplicationTests {
//存入redis //存入redis
set.add(vm.getId().toString()); set.add(vm.getId().toString());
redisCache.setCacheSet(RedisConstant.NOTIFY_VM_IDS, set); redisCache.setCacheSet(RedisConstant.NOTIFY_VM_IDS, set);
} finally {
countDownLatch.countDown(); // 完成一个线程
} }
}); });
hostMachineService.saveOrUpdateBatch(vms);
//检测物理机在线状态(检测物理机要用探测器的isTureUrl接口)
LambdaQueryWrapper<HostDetector> detectorWrapper = new LambdaQueryWrapper<>();
List<HostDetector> hostDetectors = hostDetectorService.list(detectorWrapper);
List<HostMachine> hostMachines = new ArrayList<>();
hostDetectors.forEach(detector -> {
try {
HttpRequest.get(detector.getDetectorIpAddress() + UrlConstant.DETECTOR_IS_TRUE_URL).setConnectionTimeout(5000).execute(true);
HostMachine hostMachine = hostMachineService.getById(detector.getHostMachineId());
if (Objects.nonNull(hostMachine)) {
hostMachine.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_ONLINE);
hostMachines.add(hostMachine);
} }
countDownLatch.await(); // 等待所有线程完成 } catch (Exception exception) {
hostMachineService.saveOrUpdateBatch(vms); // 保存更新后的主机状态 HostMachine hostMachine = hostMachineService.getById(detector.getHostMachineId());
executorService.shutdown(); // 关闭线程池 if (Objects.nonNull(hostMachine)) {
hostMachine.setHostMachineState(SystemConstant.HOST_MACHINE_STATE_OFFLINE);
hostMachines.add(hostMachine);
}
//判断Redis里面有没有 如果有就不需要提醒 如果没有就提醒
Set<String> set = redisCache.getCacheSet(RedisConstant.NOTIFY_HOST_IDS);
if (!CollectionUtils.isEmpty(set)) {
//如果redis里面有 说明已经发送过了未check的通知 所以不需要发送 直接返回
if (set.contains(hostMachine.getId().toString())) {
return;
} else {
//根据实例对应的提醒方式进行提醒
if (hostMachine.getNotify().equals(SystemConstant.NOTIFY_NO)) {
return;
} else {
notificationService.sendHostOfflineNotification(hostMachine.getId());
}
}
}
//存入redis
set.add(hostMachine.getId().toString());
redisCache.setCacheSet(RedisConstant.NOTIFY_HOST_IDS, set);
}
});
hostMachineService.saveOrUpdateBatch(hostMachines);
long endTime = System.nanoTime(); // 记录结束时间 long endTime = System.nanoTime(); // 记录结束时间
long elapsedTime = endTime - startTime; long elapsedTime = endTime - startTime;