diff --git a/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporter.java b/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporter.java index 6361f0edc6b5..8ad7f1d1164a 100644 --- a/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporter.java +++ b/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporter.java @@ -18,7 +18,23 @@ public interface PrometheusExporter { + /** + * Update the Prometheus metrics in text format. + * + * NOTE: capacity data is refreshed independently by {@code AlertManagerImpl}'s own + * periodic {@code CapacityChecker} timer. Do NOT force a synchronous + * {@code recalculateCapacity()} call here: it spins up a fresh thread pool per host + * and per storage pool across ALL zones on every single scrape, so with Z zones a + * single Prometheus scrape triggered Z redundant full recalculations. That extra, + * uncoordinated load compounds over time and can lead to {@code scrape_duration_seconds} + * climbing until a management-server restart. + * + * @see PrometheusExporterImpl#updateMetrics() + */ void updateMetrics(); + /** + * @return the latest Prometheus metrics refreshed by {@link #updateMetrics()}. + */ String getMetrics(); } diff --git a/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporterImpl.java b/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporterImpl.java index b49f11c77745..f737bad25ca2 100644 --- a/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporterImpl.java +++ b/plugins/integrations/prometheus/src/main/java/org/apache/cloudstack/metrics/PrometheusExporterImpl.java @@ -32,7 +32,6 @@ import org.apache.cloudstack.storage.datastore.db.ImageStoreDao; import org.apache.commons.lang3.StringUtils; -import com.cloud.alert.AlertManager; import com.cloud.api.ApiDBUtils; import com.cloud.api.query.dao.DomainJoinDao; import com.cloud.api.query.dao.StoragePoolJoinDao; @@ -126,8 +125,6 @@ public String toString() { @Inject private DomainJoinDao domainDao; @Inject - private AlertManager alertManager; - @Inject DedicatedResourceDao _dedicatedDao; @Inject private AccountDao _accountDao; @@ -497,7 +494,6 @@ public void updateMetrics() { for (final DataCenterVO dc : dcDao.listAll()) { final String zoneName = dc.getName(); final String zoneUuid = dc.getUuid(); - alertManager.recalculateCapacity(); addHostMetrics(latestMetricsItems, dc.getId(), zoneName, zoneUuid); addVMMetrics(latestMetricsItems, dc.getId(), zoneName, zoneUuid); addVolumeMetrics(latestMetricsItems, dc.getId(), zoneName, zoneUuid); diff --git a/server/src/main/java/com/cloud/alert/AlertManagerImpl.java b/server/src/main/java/com/cloud/alert/AlertManagerImpl.java index 7bf00037ee4b..9141acc994a7 100644 --- a/server/src/main/java/com/cloud/alert/AlertManagerImpl.java +++ b/server/src/main/java/com/cloud/alert/AlertManagerImpl.java @@ -161,6 +161,8 @@ public class AlertManagerImpl extends ManagerBase implements AlertManager, Confi private final ExecutorService _executor; + private ExecutorService capacityExecutorService; + protected SMTPMailSender mailSender; protected String[] recipients = null; protected String senderAddress = null; @@ -249,6 +251,9 @@ public boolean start() { @Override public boolean stop() { _timer.cancel(); + if (capacityExecutorService != null) { + capacityExecutorService.shutdown(); + } return true; } @@ -281,6 +286,24 @@ public void sendAlert(AlertType alertType, long dataCenterId, Long podId, String } } + /** + * Shared, long-lived pool for capacity recalculation, reused across every + * recalculateHostCapacities()/recalculateStorageCapacities() call instead of creating and + * tearing down a new thread pool per invocation. Repeatedly creating/shutting down pools was + * unnecessary overhead under frequent callers (e.g. the Prometheus exporter used to trigger a + * full recalculation on every scrape, see https://github.com/apache/cloudstack/issues/13586). + * Lazily created so this remains safe for callers that invoke the recalculate methods directly + * without going through configure()/start() (e.g. unit tests). + */ + private synchronized ExecutorService getCapacityExecutorService() { + if (capacityExecutorService == null || capacityExecutorService.isShutdown()) { + capacityExecutorService = Executors.newFixedThreadPool( + Math.max(1, CapacityManager.CapacityCalculateWorkers.value()), + new NamedThreadFactory("Capacity-Calculator")); + } + return capacityExecutorService; + } + /** * Recalculates the capacities of hosts, including CPU and RAM. */ @@ -290,10 +313,8 @@ protected void recalculateHostCapacities() { return; } ConcurrentHashMap> futures = new ConcurrentHashMap<>(); - ExecutorService executorService = Executors.newFixedThreadPool(Math.max(1, - Math.min(CapacityManager.CapacityCalculateWorkers.value(), hostIds.size()))); for (Long hostId : hostIds) { - futures.put(hostId, executorService.submit(() -> { + futures.put(hostId, getCapacityExecutorService().submit(() -> { final HostVO host = hostDao.findById(hostId); _capacityMgr.updateCapacityForHost(host); return null; @@ -307,7 +328,6 @@ protected void recalculateHostCapacities() { entry.getKey(), e.getMessage()), e); } } - executorService.shutdown(); } protected void recalculateStorageCapacities() { @@ -316,10 +336,8 @@ protected void recalculateStorageCapacities() { return; } ConcurrentHashMap> futures = new ConcurrentHashMap<>(); - ExecutorService executorService = Executors.newFixedThreadPool(Math.max(1, - Math.min(CapacityManager.CapacityCalculateWorkers.value(), storagePoolIds.size()))); for (Long poolId: storagePoolIds) { - futures.put(poolId, executorService.submit(() -> { + futures.put(poolId, getCapacityExecutorService().submit(() -> { Transaction.execute(new TransactionCallbackNoReturn() { @Override public void doInTransactionWithoutResult(TransactionStatus status) { @@ -343,7 +361,6 @@ public void doInTransactionWithoutResult(TransactionStatus status) { entry.getKey(), e.getMessage()), e); } } - executorService.shutdown(); } @Override diff --git a/server/src/test/java/com/cloud/alert/AlertManagerImplTest.java b/server/src/test/java/com/cloud/alert/AlertManagerImplTest.java index d34d0b5873f2..b69178b11d9d 100644 --- a/server/src/test/java/com/cloud/alert/AlertManagerImplTest.java +++ b/server/src/test/java/com/cloud/alert/AlertManagerImplTest.java @@ -17,7 +17,10 @@ package com.cloud.alert; import java.io.UnsupportedEncodingException; +import java.lang.reflect.Field; import java.util.List; +import java.util.Timer; +import java.util.concurrent.ExecutorService; import javax.mail.MessagingException; @@ -219,4 +222,114 @@ public void testRecalculateStorageCapacities() { Mockito.verify(storageManager, Mockito.times(2)).createCapacityEntry(sharedPool, Capacity.CAPACITY_TYPE_STORAGE_ALLOCATED, 10L); Mockito.verify(storageManager, Mockito.times(1)).createCapacityEntry(nonSharedPool, Capacity.CAPACITY_TYPE_LOCAL_STORAGE, 20L); } + + @Test + public void testRecalculateHostCapacitiesWithEmptyHostList() throws Exception { + Mockito.when(hostDao.listIdsByType(Host.Type.Routing)).thenReturn(List.of()); + alertManagerImplMock.recalculateHostCapacities(); + Mockito.verify(hostDao, Mockito.never()).findById(Mockito.anyLong()); + Mockito.verify(capacityManager, Mockito.never()).updateCapacityForHost(Mockito.any()); + assertNull("executor should never be created when there is nothing to submit", getCapacityExecutorService()); + } + + @Test + public void testRecalculateStorageCapacitiesWithEmptyPoolList() throws Exception { + Mockito.when(primaryDataStoreDao.listAllIds()).thenReturn(List.of()); + alertManagerImplMock.recalculateStorageCapacities(); + Mockito.verify(primaryDataStoreDao, Mockito.never()).findById(Mockito.anyLong()); + Mockito.verify(storageManager, Mockito.never()).createCapacityEntry(Mockito.any(), Mockito.anyShort(), Mockito.anyLong()); + assertNull("executor should never be created when there is nothing to submit", getCapacityExecutorService()); + } + + @Test + public void testRecalculateHostCapacitiesLogsAndContinuesOnTaskFailure() { + Mockito.when(hostDao.listIdsByType(Host.Type.Routing)).thenReturn(List.of(1L, 2L, 3L)); + HostVO host1 = Mockito.mock(HostVO.class); + HostVO host2 = Mockito.mock(HostVO.class); + HostVO host3 = Mockito.mock(HostVO.class); + Mockito.when(hostDao.findById(1L)).thenReturn(host1); + Mockito.when(hostDao.findById(2L)).thenReturn(host2); + Mockito.when(hostDao.findById(3L)).thenReturn(host3); + Mockito.doThrow(new RuntimeException("boom")).when(capacityManager).updateCapacityForHost(host2); + + alertManagerImplMock.recalculateHostCapacities(); + + Mockito.verify(capacityManager).updateCapacityForHost(host1); + Mockito.verify(capacityManager).updateCapacityForHost(host2); + Mockito.verify(capacityManager).updateCapacityForHost(host3); + Mockito.verify(alertManagerImplMock.logger).error(Mockito.anyString(), Mockito.any(Throwable.class)); + } + + @Test + public void testRecalculateHostCapacitiesReusesExecutorAcrossCalls() throws Exception { + Mockito.when(hostDao.listIdsByType(Host.Type.Routing)).thenReturn(List.of(1L)); + Mockito.when(hostDao.findById(Mockito.anyLong())).thenReturn(Mockito.mock(HostVO.class)); + + alertManagerImplMock.recalculateHostCapacities(); + ExecutorService firstExecutor = getCapacityExecutorService(); + assertNotNull(firstExecutor); + + Mockito.when(primaryDataStoreDao.listAllIds()).thenReturn(List.of(101L)); + StoragePoolVO pool = Mockito.mock(StoragePoolVO.class); + Mockito.when(primaryDataStoreDao.findById(101L)).thenReturn(pool); + alertManagerImplMock.recalculateStorageCapacities(); + + assertEquals("host and storage recalculation should share the same long-lived pool", + firstExecutor, getCapacityExecutorService()); + } + + @Test + public void testRecalculateHostCapacitiesRecreatesExecutorAfterShutdown() throws Exception { + Mockito.when(hostDao.listIdsByType(Host.Type.Routing)).thenReturn(List.of(1L)); + Mockito.when(hostDao.findById(Mockito.anyLong())).thenReturn(Mockito.mock(HostVO.class)); + + alertManagerImplMock.recalculateHostCapacities(); + ExecutorService firstExecutor = getCapacityExecutorService(); + firstExecutor.shutdown(); + + alertManagerImplMock.recalculateHostCapacities(); + ExecutorService secondExecutor = getCapacityExecutorService(); + + Assert.assertNotEquals("a shut down executor should be replaced rather than reused", firstExecutor, secondExecutor); + Assert.assertFalse(secondExecutor.isShutdown()); + } + + @Test + public void testStopShutsDownCapacityExecutorServiceWhenPresent() throws Exception { + Timer timerMock = Mockito.mock(Timer.class); + setTimer(timerMock); + Mockito.when(hostDao.listIdsByType(Host.Type.Routing)).thenReturn(List.of(1L)); + Mockito.when(hostDao.findById(Mockito.anyLong())).thenReturn(Mockito.mock(HostVO.class)); + alertManagerImplMock.recalculateHostCapacities(); + + boolean result = alertManagerImplMock.stop(); + + Assert.assertTrue(result); + Mockito.verify(timerMock).cancel(); + Assert.assertTrue(getCapacityExecutorService().isShutdown()); + } + + @Test + public void testStopDoesNotThrowWhenCapacityExecutorServiceNeverCreated() throws Exception { + Timer timerMock = Mockito.mock(Timer.class); + setTimer(timerMock); + + boolean result = alertManagerImplMock.stop(); + + Assert.assertTrue(result); + Mockito.verify(timerMock).cancel(); + assertNull(getCapacityExecutorService()); + } + + private ExecutorService getCapacityExecutorService() throws Exception { + Field field = AlertManagerImpl.class.getDeclaredField("capacityExecutorService"); + field.setAccessible(true); + return (ExecutorService) field.get(alertManagerImplMock); + } + + private void setTimer(Timer timer) throws Exception { + Field field = AlertManagerImpl.class.getDeclaredField("_timer"); + field.setAccessible(true); + field.set(alertManagerImplMock, timer); + } }