From 8cfbd4c287c2f18bb839ade9d0fe3e232cfe2249 Mon Sep 17 00:00:00 2001 From: Logic Date: Mon, 3 Aug 2026 15:45:36 +0800 Subject: [PATCH] perf: defer monitor availability detection --- .../common/constants/CommonConstants.java | 15 ++--- .../common/entity/manager/Monitor.java | 2 +- .../OldMonitorStatusWriteModelService.java | 4 +- .../service/impl/MonitorServiceImpl.java | 22 +++---- .../manager/service/MonitorServiceTest.java | 21 +++++++ ...DiscoveryExpansionSourceOwnershipTest.java | 2 +- ...OldMonitorStatusWriteModelServiceTest.java | 6 +- ...orStatusWriteModelSourceOwnershipTest.java | 4 +- .../warehouse/store/DataStorageDispatch.java | 11 ++-- .../store/DataStorageDispatchStatusTest.java | 60 +++++++++++++++++++ 10 files changed, 111 insertions(+), 36 deletions(-) create mode 100644 hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatchStatusTest.java diff --git a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/constants/CommonConstants.java b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/constants/CommonConstants.java index fe7455b1b5..16f6338c2c 100644 --- a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/constants/CommonConstants.java +++ b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/constants/CommonConstants.java @@ -57,21 +57,18 @@ public interface CommonConstants { */ byte LOGIN_FAILED_CODE = 0x05; - /** - * Monitoring status 0: Paused, 1: Up, 2: Down - */ + /** Monitoring status 0: Paused. */ byte MONITOR_PAUSED_CODE = 0x00; - /** - * Monitoring status 0: Paused, 1: Up, 2: Down - */ + /** Monitoring status 1: Up. */ byte MONITOR_UP_CODE = 0x01; - /** - * Monitoring status 0: Paused, 1: Up, 2: Down - */ + /** Monitoring status 2: Down. */ byte MONITOR_DOWN_CODE = 0x02; + /** Monitoring status 3: Scheduled and waiting for the first availability result. */ + byte MONITOR_PENDING_CODE = 0x03; + /** * scrape type static */ diff --git a/hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/entity/manager/Monitor.java b/hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/entity/manager/Monitor.java index 07dd0f8578..2a6b9c344e 100644 --- a/hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/entity/manager/Monitor.java +++ b/hertzbeat-common-spring/src/main/java/org/apache/hertzbeat/common/entity/manager/Monitor.java @@ -98,7 +98,7 @@ public class Monitor { @Size(max = 100) private String cronExpression; - @Schema(title = "Task status 0: Paused, 1: Up, 2: Down", accessMode = READ_WRITE) + @Schema(title = "Task status 0: Paused, 1: Up, 2: Down, 3: Pending", accessMode = READ_WRITE) @Min(0) @Max(4) private byte status; diff --git a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelService.java b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelService.java index 48b9462a95..b3bdcf3700 100644 --- a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelService.java +++ b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelService.java @@ -53,14 +53,14 @@ public class OldMonitorStatusWriteModelService { .collect(Collectors.toList()); } - public List findAndMarkPausedMonitorsUp(Set monitorIds) { + public List findAndMarkPausedMonitorsPending(Set monitorIds) { if (CollectionUtils.isEmpty(monitorIds)) { return List.of(); } return monitorDao.findMonitorsByIdIn(monitorIds) .stream() .filter(monitor -> monitor.getStatus() == CommonConstants.MONITOR_PAUSED_CODE) - .peek(monitor -> monitor.setStatus(CommonConstants.MONITOR_UP_CODE)) + .peek(monitor -> monitor.setStatus(CommonConstants.MONITOR_PENDING_CODE)) .collect(Collectors.toList()); } diff --git a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/impl/MonitorServiceImpl.java b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/impl/MonitorServiceImpl.java index 685e9227fa..eea211bec1 100644 --- a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/impl/MonitorServiceImpl.java +++ b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/impl/MonitorServiceImpl.java @@ -221,14 +221,12 @@ public class MonitorServiceImpl implements MonitorService { return new Configmap(param.getField(), param.getParamValue(), param.getType()); }).collect(Collectors.toList()); appDefine.setConfigmap(configmaps); + // The cyclic job owns the first availability result. Persist a distinct + // pending state instead of blocking this write on a duplicate probe or + // presenting an unobserved target as healthy/down. + monitor.setStatus(CommonConstants.MONITOR_PENDING_CODE); long jobId = collector == null ? collectJobScheduling.addAsyncCollectJob(appDefine, null) : collectJobScheduling.addAsyncCollectJob(appDefine, collector); - try { - detectMonitor(monitor, params, collector); - } catch (Exception e) { - log.warn("Monitor detection failed during addMonitor for monitor [{}]: {}", - monitor.getName(), e.getMessage()); - } try { oldMonitorCollectorBindWriteModelService.saveCollectorBind(monitorId, collector); @@ -626,7 +624,7 @@ public class MonitorServiceImpl implements MonitorService { return; } List unManagedMonitors = oldMonitorStatusWriteModelService - .findAndMarkPausedMonitorsUp(allMonitorIds); + .findAndMarkPausedMonitorsPending(allMonitorIds); if (unManagedMonitors.isEmpty()) { return; } @@ -674,12 +672,6 @@ public class MonitorServiceImpl implements MonitorService { long newJobId = collectJobScheduling.addAsyncCollectJob(appDefine, collector); monitor.setJobId(newJobId); applicationContext.publishEvent(new MonitorDeletedEvent(applicationContext, monitor.getId())); - try { - detectMonitor(monitor, params, collector); - } catch (Exception e) { - log.warn("Monitor detection failed during reapplyMonitors for monitor [{}]: {}", - monitor.getName(), e.getMessage()); - } } oldMonitorStatusWriteModelService.saveMonitorStatusChanges(unManagedMonitors); } @@ -696,6 +688,9 @@ public class MonitorServiceImpl implements MonitorService { for (AppCount item : appCounts) { AppCount appCount = appCountMap.getOrDefault(item.getApp(), new AppCount()); appCount.setApp(item.getApp()); + // Preserve the total even for transitional or future statuses that + // do not belong to the three legacy availability counters. + appCount.setSize(appCount.getSize() + item.getSize()); switch (item.getStatus()) { case CommonConstants.MONITOR_UP_CODE -> appCount.setAvailableSize(appCount.getAvailableSize() + item.getSize()); @@ -711,7 +706,6 @@ public class MonitorServiceImpl implements MonitorService { // Traverse the map obtained by statistics and convert it into a List // result set return appCountMap.values().stream().map(item -> { - item.setSize(item.getAvailableSize() + item.getUnManageSize() + item.getUnAvailableSize()); try { Job job = appService.getAppDefine(item.getApp()); item.setCategory(job.getCategory()); diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/MonitorServiceTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/MonitorServiceTest.java index fae13d1725..64642e20f7 100644 --- a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/MonitorServiceTest.java +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/MonitorServiceTest.java @@ -25,6 +25,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.Mockito.any; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.reset; import static org.mockito.Mockito.when; import java.util.ArrayList; @@ -328,6 +329,8 @@ class MonitorServiceTest { assertEquals(10, job.getDefaultInterval()); assertEquals("interval", job.getScheduleType()); assertNull(job.getCronExpression()); + assertEquals(CommonConstants.MONITOR_PENDING_CODE, monitor.getStatus()); + verify(collectJobScheduling, never()).collectSyncJobData(any(Job.class)); verify(entityIdentityResolutionService).refreshAutoMonitorBinds(monitor); } @@ -1240,6 +1243,7 @@ class MonitorServiceTest { assertEquals(10, job.getDefaultInterval()); assertEquals("cron", job.getScheduleType()); assertEquals("0 0 * * * ?", job.getCronExpression()); + monitors.forEach(monitor -> assertEquals(CommonConstants.MONITOR_PENDING_CODE, monitor.getStatus())); } @Test @@ -1278,6 +1282,23 @@ class MonitorServiceTest { assertDoesNotThrow(() -> monitorService.getAllAppMonitorsCount()); } + @Test + void getAllAppMonitorsCountIncludesPendingMonitorsInTotal() { + AppCount pending = new AppCount("test", CommonConstants.MONITOR_PENDING_CODE, 2L); + when(monitorDao.findAppsStatusCount()).thenReturn(List.of(pending)); + Job job = new Job(); + job.setMetrics(new ArrayList<>()); + when(appService.getAppDefine("test")).thenReturn(job); + + List result = monitorService.getAllAppMonitorsCount(); + + assertEquals(1, result.size()); + assertEquals(2L, result.getFirst().getSize()); + assertEquals(0L, result.getFirst().getAvailableSize()); + assertEquals(0L, result.getFirst().getUnAvailableSize()); + assertEquals(0L, result.getFirst().getUnManageSize()); + } + @Test void getMonitor() { long monitorId = 1L; diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorServiceDiscoveryExpansionSourceOwnershipTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorServiceDiscoveryExpansionSourceOwnershipTest.java index 84ac081d78..c6f410d19a 100644 --- a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorServiceDiscoveryExpansionSourceOwnershipTest.java +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorServiceDiscoveryExpansionSourceOwnershipTest.java @@ -60,7 +60,7 @@ class OldMonitorServiceDiscoveryExpansionSourceOwnershipTest { assertTrue(compactSource.indexOf("SetallMonitorIds=oldMonitorServiceDiscoveryExpansionService" + ".resolveMonitorIdsWithServiceDiscoveryChildren(ids);") < compactSource.indexOf("ListunManagedMonitors=oldMonitorStatusWriteModelService" - + ".findAndMarkPausedMonitorsUp(allMonitorIds);")); + + ".findAndMarkPausedMonitorsPending(allMonitorIds);")); assertTrue(Files.exists(OLD_MONITOR_SERVICE_DISCOVERY_EXPANSION_SERVICE)); String expansionSource = Files.readString(OLD_MONITOR_SERVICE_DISCOVERY_EXPANSION_SERVICE); diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelServiceTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelServiceTest.java index 0ec8b32a44..a45336405f 100644 --- a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelServiceTest.java +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelServiceTest.java @@ -69,17 +69,17 @@ class OldMonitorStatusWriteModelServiceTest { } @Test - void findAndMarkPausedMonitorsUpFiltersActiveRows() { + void findAndMarkPausedMonitorsPendingFiltersActiveRows() { Set monitorIds = Set.of(1L, 2L, 3L); Monitor upMonitor = Monitor.builder().id(1L).status(CommonConstants.MONITOR_UP_CODE).build(); Monitor pausedMonitor = Monitor.builder().id(2L).status(CommonConstants.MONITOR_PAUSED_CODE).build(); Monitor downMonitor = Monitor.builder().id(3L).status(CommonConstants.MONITOR_DOWN_CODE).build(); when(monitorDao.findMonitorsByIdIn(monitorIds)).thenReturn(List.of(upMonitor, pausedMonitor, downMonitor)); - List monitors = oldMonitorStatusWriteModelService.findAndMarkPausedMonitorsUp(monitorIds); + List monitors = oldMonitorStatusWriteModelService.findAndMarkPausedMonitorsPending(monitorIds); assertEquals(List.of(pausedMonitor), monitors); - assertEquals(CommonConstants.MONITOR_UP_CODE, pausedMonitor.getStatus()); + assertEquals(CommonConstants.MONITOR_PENDING_CODE, pausedMonitor.getStatus()); assertEquals(CommonConstants.MONITOR_UP_CODE, upMonitor.getStatus()); assertEquals(CommonConstants.MONITOR_DOWN_CODE, downMonitor.getStatus()); } diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelSourceOwnershipTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelSourceOwnershipTest.java index 1363f4beb8..24c9daf518 100644 --- a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelSourceOwnershipTest.java +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/OldMonitorStatusWriteModelSourceOwnershipTest.java @@ -71,7 +71,7 @@ class OldMonitorStatusWriteModelSourceOwnershipTest { assertTrue(normalizedSource.contains( "oldMonitorStatusWriteModelService.findAndMarkManagedMonitorsPaused(allMonitorIds)")); assertTrue(normalizedSource.contains( - "oldMonitorStatusWriteModelService.findAndMarkPausedMonitorsUp(allMonitorIds)")); + "oldMonitorStatusWriteModelService.findAndMarkPausedMonitorsPending(allMonitorIds)")); assertTrue(normalizedSource.contains( "oldMonitorStatusWriteModelService.saveMonitorStatusChanges(managedMonitors)")); assertTrue(normalizedSource.contains( @@ -79,7 +79,7 @@ class OldMonitorStatusWriteModelSourceOwnershipTest { String writeModelSource = Files.readString(OLD_MONITOR_STATUS_WRITE_MODEL_SERVICE); assertTrue(writeModelSource.contains("public List findAndMarkManagedMonitorsPaused(Set")); - assertTrue(writeModelSource.contains("public List findAndMarkPausedMonitorsUp(Set")); + assertTrue(writeModelSource.contains("public List findAndMarkPausedMonitorsPending(Set")); assertTrue(writeModelSource.contains("public void saveMonitorStatusChanges(List monitors)")); assertTrue(writeModelSource.contains("monitorDao.findMonitorsByIdIn(monitorIds)")); assertTrue(writeModelSource.contains("monitorDao.saveAll(monitors)")); diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatch.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatch.java index 0684cdc8cd..9563e7ab4e 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatch.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatch.java @@ -159,10 +159,13 @@ public class DataStorageDispatch { long id = metricsData.getId(); CollectRep.Code code = metricsData.getCode(); try { - String sql = "UPDATE hzb_monitor SET status = ? WHERE id = ? AND status = ?"; - int status = code == CollectRep.Code.SUCCESS ? CommonConstants.MONITOR_UP_CODE : CommonConstants.MONITOR_DOWN_CODE; - int preStatus = code == CollectRep.Code.SUCCESS ? CommonConstants.MONITOR_DOWN_CODE : CommonConstants.MONITOR_UP_CODE; - int matchedRows = jdbcTemplate.update(sql, status, id, preStatus); + String sql = "UPDATE hzb_monitor SET status = ? WHERE id = ? AND status <> ? AND status <> ?"; + byte status = code == CollectRep.Code.SUCCESS + ? CommonConstants.MONITOR_UP_CODE + : CommonConstants.MONITOR_DOWN_CODE; + // Paused monitors must remain paused. Every other non-current + // state, including Pending, converges on the first priority-0 result. + int matchedRows = jdbcTemplate.update(sql, status, id, CommonConstants.MONITOR_PAUSED_CODE, status); if (matchedRows > 0) { entityManager.getEntityManagerFactory().getCache().evict(Monitor.class, id); } diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatchStatusTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatchStatusTest.java new file mode 100644 index 0000000000..2b85e5a637 --- /dev/null +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/DataStorageDispatchStatusTest.java @@ -0,0 +1,60 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hertzbeat.warehouse.store; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +import java.util.List; +import org.apache.hertzbeat.common.constants.CommonConstants; +import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.queue.CommonDataQueue; +import org.apache.hertzbeat.plugin.runner.PluginRunner; +import org.apache.hertzbeat.warehouse.WarehouseWorkerPool; +import org.apache.hertzbeat.warehouse.store.realtime.RealTimeDataWriter; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.JdbcTemplate; + +class DataStorageDispatchStatusTest { + + @Test + void firstAvailabilityResultCanReplacePendingStatus() { + JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class); + DataStorageDispatch dispatch = new DataStorageDispatch( + mock(CommonDataQueue.class), + mock(WarehouseWorkerPool.class), + jdbcTemplate, + List.of(), + mock(RealTimeDataWriter.class), + mock(PluginRunner.class)); + CollectRep.MetricsData firstResult = CollectRep.MetricsData.newBuilder() + .setId(42L) + .setPriority(0) + .setCode(CollectRep.Code.SUCCESS) + .build(); + + dispatch.calculateMonitorStatus(firstResult); + + verify(jdbcTemplate).update( + "UPDATE hzb_monitor SET status = ? WHERE id = ? AND status <> ? AND status <> ?", + CommonConstants.MONITOR_UP_CODE, + 42L, + CommonConstants.MONITOR_PAUSED_CODE, + CommonConstants.MONITOR_UP_CODE); + } +}