perf: defer monitor availability detection

This commit is contained in:
Logic
2026-08-03 15:45:36 +08:00
parent 446faec7f6
commit 8cfbd4c287
10 changed files with 111 additions and 36 deletions
@@ -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
*/
@@ -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;
@@ -53,14 +53,14 @@ public class OldMonitorStatusWriteModelService {
.collect(Collectors.toList());
}
public List<Monitor> findAndMarkPausedMonitorsUp(Set<Long> monitorIds) {
public List<Monitor> findAndMarkPausedMonitorsPending(Set<Long> 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());
}
@@ -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<Monitor> 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<App Count>
// 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());
@@ -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<AppCount> 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;
@@ -60,7 +60,7 @@ class OldMonitorServiceDiscoveryExpansionSourceOwnershipTest {
assertTrue(compactSource.indexOf("Set<Long>allMonitorIds=oldMonitorServiceDiscoveryExpansionService"
+ ".resolveMonitorIdsWithServiceDiscoveryChildren(ids);")
< compactSource.indexOf("List<Monitor>unManagedMonitors=oldMonitorStatusWriteModelService"
+ ".findAndMarkPausedMonitorsUp(allMonitorIds);"));
+ ".findAndMarkPausedMonitorsPending(allMonitorIds);"));
assertTrue(Files.exists(OLD_MONITOR_SERVICE_DISCOVERY_EXPANSION_SERVICE));
String expansionSource = Files.readString(OLD_MONITOR_SERVICE_DISCOVERY_EXPANSION_SERVICE);
@@ -69,17 +69,17 @@ class OldMonitorStatusWriteModelServiceTest {
}
@Test
void findAndMarkPausedMonitorsUpFiltersActiveRows() {
void findAndMarkPausedMonitorsPendingFiltersActiveRows() {
Set<Long> 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<Monitor> monitors = oldMonitorStatusWriteModelService.findAndMarkPausedMonitorsUp(monitorIds);
List<Monitor> 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());
}
@@ -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<Monitor> findAndMarkManagedMonitorsPaused(Set<Long>"));
assertTrue(writeModelSource.contains("public List<Monitor> findAndMarkPausedMonitorsUp(Set<Long>"));
assertTrue(writeModelSource.contains("public List<Monitor> findAndMarkPausedMonitorsPending(Set<Long>"));
assertTrue(writeModelSource.contains("public void saveMonitorStatusChanges(List<Monitor> monitors)"));
assertTrue(writeModelSource.contains("monitorDao.findMonitorsByIdIn(monitorIds)"));
assertTrue(writeModelSource.contains("monitorDao.saveAll(monitors)"));
@@ -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);
}
@@ -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);
}
}