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 @@ -43,7 +43,7 @@ public final class ConnectionConfiguration {
/**
* Default refresh time in seconds.
*/
private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 60; //seconds
private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 10; //seconds

/**
* Default wait time before retrying to connect.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler;
import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.domain.QueueWatchDogListener;
import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

/**
* WatchDog to detect dead messages. Dead messages are messages once put on the queue, but those messages have gone. For example because
Expand Down Expand Up @@ -73,9 +74,9 @@ public void onWorkerFinished(final String messageId, final Map<String, Object> m
}

@Override
public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) {
synchronized (runningTasks) {
if (isItDead(!runningTasks.isEmpty(), numberOfMessages)) {
if (isItDead(!runningTasks.isEmpty(), queueStatus.messages())) {
LOG.info("It looks like some tasks are zombies on {} worker queue. All tasks in state running are released (running:{}).", workerQueueName,
runningTasks.size());
runningTasks.clear();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import java.util.concurrent.locks.ReentrantLock;

import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

/**
* Class to be used at startup. The Scheduler should not start before the number of messages on the queue is zero.
Expand Down Expand Up @@ -48,7 +49,7 @@ public void waitForOpen() throws InterruptedException {
}

@Override
public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) {
lock.lock();
try {
if (!open) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler;
import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.domain.QueueWatchDogListener;
import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;
import nl.aerius.taskmanager.domain.Task;
import nl.aerius.taskmanager.domain.TaskRecord;
import nl.aerius.taskmanager.domain.WorkerUpdateHandler;
Expand Down Expand Up @@ -178,13 +179,13 @@ public void reserveWorker() {
}

@Override
public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) {
synchronized (this) {
if (!firstUpdateReceived) {
initialUnaccountedWorkers = numberOfMessages;
initialUnaccountedWorkers = queueStatus.messages();
firstUpdateReceived = true;
}
updateNumberOfWorkers(numberOfWorkers);
updateNumberOfWorkers(queueStatus.consumers());
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,17 +16,17 @@
*/
package nl.aerius.taskmanager.adaptor;

import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

/**
* Interface called to observer the number of workers.
*/
public interface WorkerSizeObserver {

/**
* Gives the number of workers processes connected on the queue.
* Gives metrics on a RabbitMQ queue.
*
* @param numberOfWorkers number of number of workers processes
* @param numberOfMessages Total number of messages on the queue
* @param numberOfMessagesInProgress Number of messages being processed by the workers
* @param queueStatus RabbitMQ status metrics
*/
void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, int numberOfMessagesInProgress);
void onNumberOfWorkersUpdate(RabbitMQQueueStatus queueStatus);
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,13 +39,6 @@ public interface WorkerSizeProviderProxy {
*/
boolean removeObserver(String workerQueueName);

/**
* Triggers to get the worker queue state.
*
* @param queueName name of the worker queue
*/
void triggerWorkerQueueState(final String queueName);

/**
* Starts the worker size provider.
*
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
/*
* 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.domain;

/**
* Data record containing several collected metrics from the admin API of a single RabbitMQ queue.
*/
public record RabbitMQQueueStatus(int consumers, int messages, int unacknowledged) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,16 @@ class LoadMetric {
private final ToDoubleBiFunction<Double, Long> sumFunction;
private final Object lock = new Object();

/**
* Constructor
*
* @param countFunction Returns the value to use a count in the previous time frame. It gets 2 parameters
* - numberOfWorkers: The number of workers available in the last time period.
* - usedWorkers: The number of workers used in the last time period.
* @param sumFunction Returns the average value to report based on the 2 parameters:
* - total: the sum of all counts since the last time this metric was requested.
* - totalMeasureTime: the total time since the last time this metric was requested.
*/
public LoadMetric(final ToDoubleBiFunction<Integer, Integer> countFunction, final ToDoubleBiFunction<Double, Long> sumFunction) {
this.countFunction = countFunction;
this.sumFunction = sumFunction;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,24 +17,21 @@
package nl.aerius.taskmanager.metrics;

import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

public class RabbitMQUsageMetricsProvider implements WorkerSizeObserver, UsageMetricsProvider {

private final String workerQueueName;

private int numberOfWorkers;
private int numberOfMessages;
private int numberOfMessagesInProgress;
private RabbitMQQueueStatus queueStatus = new RabbitMQQueueStatus(0,0,0);

public RabbitMQUsageMetricsProvider(final String workerQueueName) {
this.workerQueueName = workerQueueName;
}

@Override
public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
this.numberOfWorkers = numberOfWorkers;
this.numberOfMessages = numberOfMessages;
this.numberOfMessagesInProgress = numberOfMessagesInProgress;
public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) {
this.queueStatus = queueStatus;
}

@Override
Expand All @@ -44,22 +41,22 @@ public String getWorkerQueueName() {

@Override
public int getNumberOfWorkers() {
return numberOfWorkers;
return queueStatus.consumers();
}

@Override
public int getNumberOfUsedWorkers() {
return numberOfMessagesInProgress;
return queueStatus.unacknowledged();
}

@Override
public int getNumberOfFreeWorkers() {
return Math.max(0, numberOfWorkers - numberOfMessagesInProgress);
return Math.max(0, getNumberOfWorkers() - getNumberOfUsedWorkers());
}

@Override
public int getNumberOfWaiting() {
return Math.max(0, numberOfMessages - numberOfMessagesInProgress);
return Math.max(0, queueStatus.messages() - getNumberOfUsedWorkers());
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler;
import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.domain.QueueWatchDogListener;
import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

/**
* This class provides the input for the {@link TaskManagerUsageMetricsProvider}. It will register updates on the amount of worker/workers from
Expand Down Expand Up @@ -57,12 +58,14 @@ public void onWorkerFinished(final String messageId, final Map<String, Object> m
}

@Override
public synchronized void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
this.numberOfWorkers = numberOfWorkers;
if (!startupGuard.isOpen() && numberOfMessages > 0) {
public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) {
this.numberOfWorkers = queueStatus.consumers();
final int messages = queueStatus.messages();

if (!startupGuard.isOpen() && messages > 0) {
LOG.info("Queue {} will be started with {} messages already on the queue.", taskManagerUsageMetricsProvider.getWorkerQueueName(),
numberOfMessages);
taskManagerUsageMetricsProvider.register(numberOfMessages, numberOfWorkers);
messages);
taskManagerUsageMetricsProvider.register(messages, numberOfWorkers);
} else {
taskManagerUsageMetricsProvider.register(0, numberOfWorkers);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,19 +20,19 @@

/**
* Class that wraps all Task Manager usage metrics reporters. TaskSchedulerBuckets should add the metric providers to this class.
* Each of those providers added is for a specific worker queue. The {@link UsageMetricsWrapper} will manager metrics per worker queue.
* Each of those providers added is for a specific worker queue. The {@link WorkerUsageMetricsWrapper} will manager metrics per worker queue.
*/
public class TaskManagerUsageMetricsWrapper {

private final UsageMetricsWrapper rabbitMQUsageMetrics;
private final UsageMetricsWrapper workerPoolUsageMetrics;
private final UsageMetricsWrapper taskManagerUsageMetrics;
private final WorkerUsageMetricsWrapper rabbitMQUsageMetrics;
private final WorkerUsageMetricsWrapper workerPoolUsageMetrics;
private final WorkerUsageMetricsWrapper taskManagerUsageMetrics;
private final UsageMetricsReporter loadUsageMetricsReporter;

public TaskManagerUsageMetricsWrapper(final Meter meter) {
rabbitMQUsageMetrics = new UsageMetricsWrapper(meter, "aer.rabbitmq", true);
workerPoolUsageMetrics = new UsageMetricsWrapper(meter, "aer.taskmanager.workerpool", false);
taskManagerUsageMetrics = new UsageMetricsWrapper(meter, "aer.taskmanager", false);
rabbitMQUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.rabbitmq", true);
workerPoolUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.taskmanager.workerpool", false);
taskManagerUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.taskmanager", false);
loadUsageMetricsReporter = new UsageMetricsReporter(meter, "aer.taskmanager.work.load", "Report average worker load");
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,15 +17,12 @@
package nl.aerius.taskmanager.metrics;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.DoubleSupplier;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.Meter;
import io.opentelemetry.api.metrics.ObservableDoubleGauge;
Expand All @@ -37,11 +34,9 @@
*/
class UsageMetricsReporter {

private static final Logger LOG = LoggerFactory.getLogger(UsageMetricsReporter.class);

private record UsageMetric(DoubleSupplier metricSupplier, Attributes attributes) {}

private final Map<String, List<UsageMetric>> metricsMap = new HashMap<>();
private final Map<String, List<UsageMetric>> metricsMap = new ConcurrentHashMap<>();
private final ObservableDoubleGauge gauge;

public UsageMetricsReporter(final Meter meter, final String metricName, final String description) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,12 @@
* The metrics supported are limit (i.e. number of workers available).
*
*/
class UsageMetricsWrapper {
class WorkerUsageMetricsWrapper {
private final boolean hasWaiting;
private final UsageMetricsReporter limitReporter;
private final UsageMetricsReporter usageReporter;

public UsageMetricsWrapper(final Meter meter, final String metricPrefix, final boolean hasWaiting) {
public WorkerUsageMetricsWrapper(final Meter meter, final String metricPrefix, final boolean hasWaiting) {
this.hasWaiting = hasWaiting;
limitReporter = new UsageMetricsReporter(meter, metricPrefix + ".worker.limit", "Report nunber of workers available");
usageReporter = new UsageMetricsReporter(meter, metricPrefix + ".worker.usage", "Report worker usage");
Expand Down
Loading
Loading