Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ class LoadMetric {
private int numberOfWorkers;
private final ToDoubleBiFunction<Integer, Integer> countFunction;
private final ToDoubleBiFunction<Double, Long> sumFunction;
private final Object lock = new Object();

public LoadMetric(final ToDoubleBiFunction<Integer, Integer> countFunction, final ToDoubleBiFunction<Double, Long> sumFunction) {
this.countFunction = countFunction;
Expand All @@ -62,37 +63,51 @@ public LoadMetric(final ToDoubleBiFunction<Integer, Integer> 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) {
usedWorkers = 0;
process();
}
}

/**
* Calculates average duration over the last time frame since this method was called. Internals are reset in this method to a new measure point.
*
* @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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<String> dispatchedTasks = ConcurrentHashMap.newKeySet();

private int numberOfWorkers;

Expand All @@ -53,19 +48,12 @@ public TaskManagerMetricsRegister(final TaskManagerUsageMetricsProvider taskMana

@Override
public void onWorkDispatched(final String messageId, final Map<String, Object> messageMetaData) {
synchronized (dispatchedTasks) {
dispatchedTasks.add(messageId);
taskManagerUsageMetricsProvider.register(1, numberOfWorkers);
}
taskManagerUsageMetricsProvider.register(1, numberOfWorkers);
}

@Override
public void onWorkerFinished(final String messageId, final Map<String, Object> messageMetaData) {
synchronized (dispatchedTasks) {
if (dispatchedTasks.remove(messageId)) {
taskManagerUsageMetricsProvider.register(-1, numberOfWorkers);
}
}
taskManagerUsageMetricsProvider.register(-1, numberOfWorkers);
}

@Override
Expand All @@ -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();
}
}
Original file line number Diff line number Diff line change
@@ -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<Integer, Integer> countFunction;
@Captor ArgumentCaptor<Integer> countUsedWorkersCaptor;
@Captor ArgumentCaptor<Integer> 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 get original number of workers");
assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get 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 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());
// 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");
Comment thread
tom-h42 marked this conversation as resolved.
loadMetric.reset();
verify(countFunction, times(3)).applyAsDouble(any(), countUsedWorkersCaptor.capture());
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 previous value of 0 used workers");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ public class TaskManagerMetricsRegisterTest {

private @Mock TaskManagerUsageMetricsProvider taskManagerUsageMetricsProvider;

private @Captor ArgumentCaptor<Integer> taskManagerUsagerMetricsProviderCaptor;
private @Captor ArgumentCaptor<Integer> startUpNrOfMessagesCaptor;
private @Captor ArgumentCaptor<Integer> lastNrOfMessagesCaptor;

private StartupGuard startupGuard;
private TaskManagerMetricsRegister register;
Expand All @@ -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
Expand All @@ -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<Integer> 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);
}

Expand Down
Loading