From 34fca6735152c98f0bd8b4ee4845d80c0dddd4dc Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Thu, 3 Sep 2026 14:39:43 +0200 Subject: [PATCH 1/3] AER-4645 Fix to update load metric also on a regular basis Current implementation only updated load when task dispatched or finished. With long running jobs this could mean it takes some time to update. But since load is used to scale it could be the load value triggered scaling, when scaled it would be expected the load to decrease (more workers for same work). But when load is not updated the load stays the same and system might still think it needs scaling up, and will scale up even more. With this change it will use the regular called update that was only used to track the initial state of the queue. It now uses a delta of 0 as there are no new number of messages. In the LoadMetric it will than only perform an update if the number of workers changed. In other situation there is no reason to update. LoadMetric changed a bit because process computes the values at the moment called but it does that by calling with last known values and delta 0. But with new guard it would result in update not being performed. Therefore put update in separate method. Removed dispatchedTasks from TaskManagerUsageMetricsProvider as it was to cover for tasks that were still on the queue when the taskmanager would be (re)started. But with startupGuard these messages are accounted for and therefore no need to filter those out is needed. --- .../taskmanager/metrics/LoadMetric.java | 52 ++++++---- .../metrics/TaskManagerMetricsRegister.java | 23 +---- .../taskmanager/metrics/LoadMetricTest.java | 96 +++++++++++++++++++ .../TaskManagerMetricsRegisterTest.java | 27 ++++-- 4 files changed, 153 insertions(+), 45 deletions(-) create mode 100644 source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java index 3f32664..63104ea 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java @@ -50,6 +50,7 @@ class LoadMetric { private int numberOfWorkers; private final ToDoubleBiFunction countFunction; private final ToDoubleBiFunction sumFunction; + private final Object lock = new Object(); public LoadMetric(final ToDoubleBiFunction countFunction, final ToDoubleBiFunction sumFunction) { this.countFunction = countFunction; @@ -62,23 +63,22 @@ public LoadMetric(final ToDoubleBiFunction countFunction, fina * @param deltaUsedWorkers number of jobs on the workers being added or subtracted. * @param numberOfWorkers Number of available workers */ - public synchronized void register(final int deltaUsedWorkers, final int numberOfWorkers) { - this.numberOfWorkers = numberOfWorkers; - final long newLast = System.currentTimeMillis(); - final long delta = newLast - last; - - total += delta * countFunction.applyAsDouble(numberOfWorkers, usedWorkers); - totalMeasureTime += delta; - last = newLast; - usedWorkers = Math.max(0, usedWorkers + deltaUsedWorkers); + public void register(final int deltaUsedWorkers, final int numberOfWorkers) { + synchronized (lock) { + if (deltaUsedWorkers != 0 || this.numberOfWorkers != numberOfWorkers) { + updateState(deltaUsedWorkers, numberOfWorkers); + } + } } + /** * Resets the metric state. Sets running workers to 0, and resets the average load time by calling process. */ - public synchronized void reset() { - usedWorkers = 0; - process(); + public void reset() { + synchronized(lock) { + updateState(-usedWorkers, numberOfWorkers); + } } /** @@ -86,13 +86,27 @@ public synchronized void reset() { * * @return Average load of the workers since the last time this method was called */ - public synchronized double process() { - // Call register here to set the end time this moment. This will calculate workers running up till now as being active. - register(0, numberOfWorkers); - final double averageTotal = totalMeasureTime > 0 ? sumFunction.applyAsDouble(total, totalMeasureTime) : 0; + public double process() { + synchronized(lock) { + // Call register here to set the end time this moment. This will calculate workers running up till now as being active. + updateState(0, numberOfWorkers); + final double averageTotal = totalMeasureTime > 0 ? sumFunction.applyAsDouble(total, totalMeasureTime) : 0; - totalMeasureTime = 0; - total = 0; - return averageTotal; + totalMeasureTime = 0; + total = 0; + return averageTotal; + } + } + + private void updateState(final int deltaUsedWorkers, final int updatedNumberOfWorkers) { + final long newLast = System.currentTimeMillis(); + final long delta = newLast - last; + + // First calculate the total for passed time period with the worker values up to this time. + total += delta * countFunction.applyAsDouble(numberOfWorkers, usedWorkers); + totalMeasureTime += delta; + last = newLast; + numberOfWorkers = updatedNumberOfWorkers; + usedWorkers = Math.max(0, usedWorkers + deltaUsedWorkers); } } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java index 56e3132..00e828f 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java @@ -17,8 +17,6 @@ package nl.aerius.taskmanager.metrics; import java.util.Map; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,9 +38,6 @@ public class TaskManagerMetricsRegister implements WorkerProducerHandler, Worker private final TaskManagerUsageMetricsProvider taskManagerUsageMetricsProvider; private final StartupGuard startupGuard; - // Keep track of dispatched tasks, because when taskmanager restarts it should not register tasks already on the queue - // as it doesn't have any metrics on it anymore. - private final Set dispatchedTasks = ConcurrentHashMap.newKeySet(); private int numberOfWorkers; @@ -53,19 +48,12 @@ public TaskManagerMetricsRegister(final TaskManagerUsageMetricsProvider taskMana @Override public void onWorkDispatched(final String messageId, final Map messageMetaData) { - synchronized (dispatchedTasks) { - dispatchedTasks.add(messageId); - taskManagerUsageMetricsProvider.register(1, numberOfWorkers); - } + taskManagerUsageMetricsProvider.register(1, numberOfWorkers); } @Override public void onWorkerFinished(final String messageId, final Map messageMetaData) { - synchronized (dispatchedTasks) { - if (dispatchedTasks.remove(messageId)) { - taskManagerUsageMetricsProvider.register(-1, numberOfWorkers); - } - } + taskManagerUsageMetricsProvider.register(-1, numberOfWorkers); } @Override @@ -75,14 +63,13 @@ public synchronized void onNumberOfWorkersUpdate(final int numberOfWorkers, fina LOG.info("Queue {} will be started with {} messages already on the queue.", taskManagerUsageMetricsProvider.getWorkerQueueName(), numberOfMessages); taskManagerUsageMetricsProvider.register(numberOfMessages, numberOfWorkers); + } else { + taskManagerUsageMetricsProvider.register(0, numberOfWorkers); } } @Override public void reset() { - synchronized (dispatchedTasks) { - dispatchedTasks.clear(); - taskManagerUsageMetricsProvider.reset(); - } + taskManagerUsageMetricsProvider.reset(); } } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java new file mode 100644 index 0000000..2d8bbdc --- /dev/null +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java @@ -0,0 +1,96 @@ +/* + * Copyright (c) Contributors to the project + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see http://www.gnu.org/licenses/. + */ +package nl.aerius.taskmanager.metrics; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import java.util.function.ToDoubleBiFunction; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +/** + * Test class for {@link LoadMetric}. + */ +@ExtendWith(MockitoExtension.class) +class LoadMetricTest { + + @Mock ToDoubleBiFunction countFunction; + @Captor ArgumentCaptor countUsedWorkersCaptor; + @Captor ArgumentCaptor countWorkersCaptor; + + private LoadMetric loadMetric; + + @BeforeEach + void beforeEach() { + doAnswer(a -> ((Integer) a.getArgument(0)).doubleValue()).when(countFunction).applyAsDouble(any(), any()); + + loadMetric = new LoadMetric(countFunction, (total, time) -> total); + } + + /** + * Register + */ + @Test + void testProcess() throws InterruptedException { + loadMetric.register(5, 10); + loadMetric.process(); + verify(countFunction, times(2)).applyAsDouble(countWorkersCaptor.capture(),countUsedWorkersCaptor.capture()); + assertEquals(10, countWorkersCaptor.getValue(), "Should got original number of workers"); + assertEquals(5, countUsedWorkersCaptor.getValue(), "Should got original number of used workers"); + } + + @Test + void testWorkersChanged() throws InterruptedException { + loadMetric.register(0, 10); + loadMetric.register(0, 10); + // 2nd call to register should not trigger updating internal state. + verify(countFunction, times(1)).applyAsDouble(countWorkersCaptor.capture(), countUsedWorkersCaptor.capture()); + assertEquals(0, countWorkersCaptor.getValue(), "Should got inital number of workers, which was 0"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "Should got inital number of used workers, which was 0"); + loadMetric.register(0, 5); + // call with changed number of workers should trigger update state. + verify(countFunction, times(2)).applyAsDouble(any(), any()); + // process call should also trigger update. + loadMetric.process(); + verify(countFunction, times(3)).applyAsDouble(any(), any()); + } + + @Test + void testReset() throws InterruptedException { + // call 2 times. because first time used workers is initialized. + loadMetric.register(5, 10); + loadMetric.register(6, 10); + verify(countFunction, times(2)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get the number of used workers of the first call"); + loadMetric.reset(); + verify(countFunction, times(3)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(11, countUsedWorkersCaptor.getValue(), "reset should set the number of used workers of the last register call"); + loadMetric.register(8, 10); + verify(countFunction, times(4)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(0, countUsedWorkersCaptor.getValue(), "reset should set used workers to 0"); + } +} diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java index af8805a..2608ca9 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java @@ -45,7 +45,8 @@ public class TaskManagerMetricsRegisterTest { private @Mock TaskManagerUsageMetricsProvider taskManagerUsageMetricsProvider; - private @Captor ArgumentCaptor taskManagerUsagerMetricsProviderCaptor; + private @Captor ArgumentCaptor startUpNrOfMessagesCaptor; + private @Captor ArgumentCaptor lastNrOfMessagesCaptor; private StartupGuard startupGuard; private TaskManagerMetricsRegister register; @@ -59,18 +60,29 @@ void beforeEach() { @Test void testOnWorkDispatched() { startUp(10, 0); + // Should have called register only once, but with 0 delta + verifyTaskManagerUsageMetricsProvider(1, 0, startUpNrOfMessagesCaptor); register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkDispatched("2", createMap(QUEUE_2, 200L)); - verifyTaskManagerUsageMetricsProvider(2, 2); + register.onNumberOfWorkersUpdate(10, 2, 2); + // Should have called register 4 times (startup, 2 for dispatch and 1 for update. + // But total delta should be +2 for the 2 dispatched messages. + verifyTaskManagerUsageMetricsProvider(4, 2, lastNrOfMessagesCaptor); } @Test void testOnWorkerFinished() { - startUp(2, 0); + startUp(2, 2); + // Should have called register only once, but with 2 as delta as it should match total number of messages. + verifyTaskManagerUsageMetricsProvider(1, 2, startUpNrOfMessagesCaptor); register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("2", createMap(QUEUE_2, 200L)); - verifyTaskManagerUsageMetricsProvider(2, 0); + register.onNumberOfWorkersUpdate(10, 2, 2); + + // Should have called register 5 times (startup, 1 for dispatch, 2 for finished and 1 for update. + // But total delta should be 0 for 2 startup + 1 dispatch - 2 for finish.. + verifyTaskManagerUsageMetricsProvider(5, 1, lastNrOfMessagesCaptor); } @Test @@ -85,10 +97,9 @@ void testReset() throws InterruptedException { verify(taskManagerUsageMetricsProvider, times(1)).reset(); } - private void verifyTaskManagerUsageMetricsProvider(final int times, final int sum) { - verify(taskManagerUsageMetricsProvider, times(times)).register(taskManagerUsagerMetricsProviderCaptor.capture(), - anyInt()); - assertEquals(sum, taskManagerUsagerMetricsProviderCaptor.getAllValues().stream().mapToInt(Integer::intValue).sum(), + private void verifyTaskManagerUsageMetricsProvider(final int times, final int sum, final ArgumentCaptor captor) { + verify(taskManagerUsageMetricsProvider, times(times)).register(captor.capture(), anyInt()); + assertEquals(sum, captor.getAllValues().stream().mapToInt(Integer::intValue).sum(), "Should have registered a total sum of " + sum); } From e17e8586c8906c5a32064317845905920cac3ab9 Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Fri, 11 Sep 2026 12:01:27 +0200 Subject: [PATCH 2/3] Review comment --- .../main/java/nl/aerius/taskmanager/metrics/LoadMetric.java | 5 +++-- .../java/nl/aerius/taskmanager/metrics/LoadMetricTest.java | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java index 63104ea..1bca565 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java @@ -63,7 +63,7 @@ public LoadMetric(final ToDoubleBiFunction countFunction, fina * @param deltaUsedWorkers number of jobs on the workers being added or subtracted. * @param numberOfWorkers Number of available workers */ - public void register(final int deltaUsedWorkers, final int numberOfWorkers) { + public void register(final int deltaUsedWorkers, final int numberOfWorkers) { synchronized (lock) { if (deltaUsedWorkers != 0 || this.numberOfWorkers != numberOfWorkers) { updateState(deltaUsedWorkers, numberOfWorkers); @@ -77,7 +77,8 @@ public void register(final int deltaUsedWorkers, final int numberOfWorkers) { */ public void reset() { synchronized(lock) { - updateState(-usedWorkers, numberOfWorkers); + usedWorkers = 0; + process(); } } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java index 2d8bbdc..141958d 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java @@ -88,9 +88,9 @@ void testReset() throws InterruptedException { assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get the number of used workers of the first call"); loadMetric.reset(); verify(countFunction, times(3)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); - assertEquals(11, countUsedWorkersCaptor.getValue(), "reset should set the number of used workers of the last register call"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "Reset has set used number to 0, so that is at is expected here."); loadMetric.register(8, 10); verify(countFunction, times(4)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); - assertEquals(0, countUsedWorkersCaptor.getValue(), "reset should set used workers to 0"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "This call should get prevoious value of 0 used workers"); } } From 7d1ba3e4e8b882cc9b8c1c3d88af8762ef74b657 Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Tue, 15 Sep 2026 09:47:43 +0200 Subject: [PATCH 3/3] Review comment --- .../nl/aerius/taskmanager/metrics/LoadMetricTest.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java index 141958d..c24c625 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java @@ -59,8 +59,8 @@ void testProcess() throws InterruptedException { loadMetric.register(5, 10); loadMetric.process(); verify(countFunction, times(2)).applyAsDouble(countWorkersCaptor.capture(),countUsedWorkersCaptor.capture()); - assertEquals(10, countWorkersCaptor.getValue(), "Should got original number of workers"); - assertEquals(5, countUsedWorkersCaptor.getValue(), "Should got original number of used workers"); + assertEquals(10, countWorkersCaptor.getValue(), "Should get original number of workers"); + assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get original number of used workers"); } @Test @@ -69,8 +69,8 @@ void testWorkersChanged() throws InterruptedException { loadMetric.register(0, 10); // 2nd call to register should not trigger updating internal state. verify(countFunction, times(1)).applyAsDouble(countWorkersCaptor.capture(), countUsedWorkersCaptor.capture()); - assertEquals(0, countWorkersCaptor.getValue(), "Should got inital number of workers, which was 0"); - assertEquals(0, countUsedWorkersCaptor.getValue(), "Should got inital number of used workers, which was 0"); + assertEquals(0, countWorkersCaptor.getValue(), "Should get inital number of workers, which was 0"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "Should get inital number of used workers, which was 0"); loadMetric.register(0, 5); // call with changed number of workers should trigger update state. verify(countFunction, times(2)).applyAsDouble(any(), any()); @@ -91,6 +91,6 @@ void testReset() throws InterruptedException { assertEquals(0, countUsedWorkersCaptor.getValue(), "Reset has set used number to 0, so that is at is expected here."); loadMetric.register(8, 10); verify(countFunction, times(4)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); - assertEquals(0, countUsedWorkersCaptor.getValue(), "This call should get prevoious value of 0 used workers"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "This call should get previous value of 0 used workers"); } }