Skip to content
Open
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
18 changes: 10 additions & 8 deletions doc/telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,9 +82,9 @@ Therefore `.limit` gives the total number of workers available.
And `.usage` gives the metrics about how the workers are used.
The usage metrics have an attribute `state` that identifies if the metric is for `used` or `free` amount of workers.
Additional `aer.rabbitmq.worker.usage` records a third state `waiting`.
And `aer.taskmanager.client.queue.usage` has states `used` and `waiting` that provides information on tasks send to the TaskManager.
This shows the number of tasks on the worker queue that are not yet picked up by any worker.

The `aer.taskmanager.queue` metrics has state attributes `used` and `waiting` that provides information on tasks picked up by the TaskManager.
Waiting tasks are only interesting for queues that use `eagerFetch`, because in that configuration all tasks are always picked up from the queue.
To know the number of waiting tasks for other queues use the metric `aer.rabbitmq.client.queue`.

| metric name | type | description |
|-------------------------------------------------------|-----------|----------------------------------------------------------------------------|
Expand All @@ -97,13 +97,14 @@ This shows the number of tasks on the worker queue that are not yet picked up by
| `aer.taskmanager.worker.usage`<sup>2</sup> | gauge | Weighted usage of workers based on tasks send to workers. |
| `aer.taskmanager.workerpool.worker.limit`<sup>1</sup> | gauge | TaskManager internal total number of workers. |
| `aer.taskmanager.workerpool.worker.usage`<sup>2</sup> | gauge | TaskManager internal usage of workers. |
| `aer.taskmanager.client.queue.usage`<sup>2</sup> | gauge | TaskManager internal metrics on client queue usage. |
| `aer.taskmanager.queue`<sup>4</sup> | gauge | TaskManager internal metrics on client queue usage. |
| `aer.rabbitmq.worker.limit`<sup>1</sup> | gauge | Total number of workers available as reported by RabbitMQ |
| `aer.rabbitmq.worker.usage`<sup>2</sup> | gauge | Usage o the workers based on the messages on the RabbitMQ worker queue. |
| `aer.rabbitmq.worker.usage`<sup>2</sup> | gauge | Usage of the workers based on the messages on the RabbitMQ worker queue. |
| `aer.rabbitmq.client.queue`<sup>3</sup> | gauge | Number of messages on per client queue as reported by RabbitMQ. |
| `aer.taskmanager.dispatched`<sup>1</sup> | histogram | The number of tasks dispatched. |
| `aer.taskmanager.dispatched.wait`<sup>1</sup> | histogram | The average wait time of tasks dispatched. |
| `aer.taskmanager.dispatched.queue`<sup>3</sup> | histogram | The number of tasks dispatched per client queue. |
| `aer.taskmanager.dispatched.queue.wait`<sup>3</sup> | histogram | The average wait time of tasks dispatched per client queue. |
| `aer.taskmanager.dispatched.queue`<sup>4</sup> | histogram | The number of tasks dispatched per client queue. |
| `aer.taskmanager.dispatched.queue.wait`<sup>4</sup> | histogram | The average wait time of tasks dispatched per client queue. |

Basically there are 3 metric groups that report similar information.
First the `aer.taskmanager.worker.*` metrics are a weighted value based on when tasks are send to the workers.
Expand All @@ -120,9 +121,10 @@ The workers have different attributes to distinguish specific metrics.
* <sup>1</sup> have attribute `worker_type`.
* <sup>2</sup> have attribute `worker_type` and `state`. `state` can have the value `used`, `free` or `waiting`.
* <sup>3</sup> have attribute `worker_type` and `queue_name`.
* <sup>4</sup> have attribute `worker_type`, `queue_name` and `state`. `state` can have the value `used`, `free` or `waiting`.

`worker_type` is the type of worker, e.g. `ops`.
`queue_name` is the originating queue the task initially was put on, e.g. `...calculator_ui_small`.
`queue_name` is the originating queue the task initially was put on, e.g. `calculator_ui_small`.

> [!NOTE]
> The metrics in the TaskManager operate on a time frame of 1 minute.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import nl.aerius.taskmanager.domain.TaskSchedule;
import nl.aerius.taskmanager.metrics.OpenTelemetryMetrics;
import nl.aerius.taskmanager.metrics.PerformanceMetricsReporter;
import nl.aerius.taskmanager.metrics.RabbitMQClientQueueReporter;
import nl.aerius.taskmanager.metrics.RabbitMQUsageMetricsProvider;
import nl.aerius.taskmanager.metrics.TaskManagerMetricsRegister;
import nl.aerius.taskmanager.metrics.TaskManagerUsageMetricsProvider;
Expand Down Expand Up @@ -167,6 +168,7 @@ public TaskScheduleBucket(final QueueConfig queueConfig) throws InterruptedExcep
workerSizeObserverProxy.addObserver(workerQueueName, taskManagerMetricsRegister);
workerSizeObserverProxy.addObserver(workerQueueName, workerPool);
workerSizeObserverProxy.addObserver(workerQueueName, watchDog);
workerSizeObserverProxy.addClientObserver(workerQueueName, new RabbitMQClientQueueReporter(OpenTelemetryMetrics.METER, workerQueueName));
// startup Guard should be the last observer added as it will unlock the task dispatcher
workerSizeObserverProxy.addObserver(workerQueueName, startupGuard);

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
/*
* 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.adaptor;

import nl.aerius.taskmanager.domain.RabbitMQQueueStatus;

/**
* Observer that listens to client queue updates.
*/
public interface ClientQueueObserver {

/**
* Update for the given client queue.
*
* @param clientQueueName name of the client queue
* @param status queue metrics
*/
void onClientQueueUpdate(String clientQueueName, RabbitMQQueueStatus status);

/**
* Returns true if this observer should receive updates for the given client queue.
*
* @param clientQueueName client queue to check
* @return true if should receive updates
*/
boolean filter(String clientQueueName);
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@ public interface WorkerSizeProviderProxy {
*/
void addObserver(String workerQueueName, WorkerSizeObserver workerSizeObserver);

/**
* Add a ClientQueueObserver.
*
* @param workerQueueName name of the worker queue the client queues are related too
* @param observer observer to be informed
*/
public void addClientObserver(final String workerQueueName, final ClientQueueObserver observer);

/**
* Removes the observer for the worker queue.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ public static Attributes queueAttributes(final String workerQueueName, final Str
.build();
}

private static String onlyLastPart(final String name) {
public static String onlyLastPart(final String name) {
final int lastIndex = name.lastIndexOf('.');

return (lastIndex < 0 ? name : name.substring(lastIndex + 1)).toLowerCase(Locale.ROOT);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
/*
* 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 java.util.HashMap;
import java.util.Map;

import io.opentelemetry.api.metrics.Meter;

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

/**
* Reports RabbitMQ client queue number of messages on the queue.
*
* The report registers the metrics using the last part of the worker queue name. Typical worker queue name is 'aerius.worker.some_name'.
* The last part of the worker queue is used as attribute when reporting the metric.
* All client queues related to this worker queue will be reported on.
* The specific client metric is reported with the last part of the client queue name as attribute value.
* The metrics are reported as gauge metric to the metric 'aer.rabbitmq.client_queue'.
*/
public class RabbitMQClientQueueReporter implements ClientQueueObserver {

private static final String CLIENT_QUEUE_METRIC = "aer.rabbitmq.client.queue";
private final String workerQueue;
private final UsageMetricsReporter clientQueueReporter;
private final Map<String, Double> queueCounters = new HashMap<>();

public RabbitMQClientQueueReporter(final Meter meter, final String workerQueueName) {
this.workerQueue = OpenTelemetryMetrics.onlyLastPart(workerQueueName);
clientQueueReporter = new UsageMetricsReporter(meter, CLIENT_QUEUE_METRIC, "Number of messages on the RabbitMQ client queues");
}

@Override
public void onClientQueueUpdate(final String clientQueueName, final RabbitMQQueueStatus value) {
if (!queueCounters.containsKey(clientQueueName)) {
clientQueueReporter.addMetrics(workerQueue, () -> queueCounters.get(clientQueueName),
OpenTelemetryMetrics.queueAttributes(workerQueue, clientQueueName));
}
queueCounters.put(clientQueueName, Double.valueOf(value.messages()));
}

@Override
public boolean filter(final String queueName) {
return queueName.contains(workerQueue) && !queueName.contains("worker");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -47,8 +47,8 @@ public UsageMetricsReporter(final Meter meter, final String metricName, final St
}

private void recordMetrics(final ObservableDoubleMeasurement measurement) {
for (final Entry<String, List<UsageMetric>> workerMetrics : metricsMap.entrySet()) {
for (final UsageMetric metric : workerMetrics.getValue()) {
for (final Entry<String, List<UsageMetric>> entry : metricsMap.entrySet()) {
for (final UsageMetric metric : entry.getValue()) {

measurement.record(metric.metricSupplier().getAsDouble(), metric.attributes());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -110,25 +110,25 @@ public void close() {
/**
* Retrieves the status for all queues from the RabbitMQ admin api.
*/
public Map<String, RabbitMQQueueStatus> getWorkerQueueStates() {
public Map<String, RabbitMQQueueStatus> getQueueStates() {
final Map<String, RabbitMQQueueStatus> queueStates = new HashMap<>();

try {
final JsonNode jsonObject = getJsonResultFromApi("/api/queues");

if (jsonObject == null) {
LOG.error("Queue configuration from RabbitMQ admin json get call returned null.");
} if (jsonObject instanceof final ArrayNode array) {
final Map<String, RabbitMQQueueStatus> queueStates = new HashMap<>();
if (jsonObject instanceof final ArrayNode array) {

for (int i = 0; i < array.size(); i++) {
final JsonNode jsonNode = array.get(i);
queueStates.put(getJsonString(jsonNode, "name"), getQueueStatus(jsonNode));
}
return queueStates;
} else {
LOG.error("Queue configuration from RabbitMQ admin json returned an unexpected value.");
}
} catch (final URISyntaxException | IOException e) {
LOG.info("Error getting RabbitMQ status from admin api: {}", e.getMessage());
}
return new HashMap<>();
return queueStates;
}

private static RabbitMQQueueStatus getQueueStatus(final JsonNode jsonNode) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import nl.aerius.taskmanager.adaptor.ClientQueueObserver;
import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerMetrics;
import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;
import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy;
Expand Down Expand Up @@ -56,6 +57,7 @@ public class RabbitMQWorkerSizeProvider implements WorkerSizeProviderProxy {
private final long refreshRateSeconds;

private final Map<String, WorkerSizeObserverComposite> observers = new HashMap<>();
private final Map<String, ClientQueueObserver> clientObservers = new HashMap<>();
private final RabbitMQQueueMonitor monitor;
private boolean running;

Expand All @@ -72,17 +74,22 @@ public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService
}

@Override
public void addObserver(final String queueName, final WorkerSizeObserver observer) {
observers.computeIfAbsent(queueName, k -> new WorkerSizeObserverComposite()).add(observer);
public void addObserver(final String workerQueueName, final WorkerSizeObserver observer) {
observers.computeIfAbsent(workerQueueName, k -> new WorkerSizeObserverComposite()).add(observer);
if (observer instanceof WorkerMetrics) {
eventProducer.addMetrics(queueName, (WorkerMetrics) observer);
eventProducer.addMetrics(workerQueueName, (WorkerMetrics) observer);
}
}

@Override
public void addClientObserver(final String workerQueueName, final ClientQueueObserver observer) {
clientObservers.put(workerQueueName, observer);
}

@Override
public boolean removeObserver(final String queueName) {
eventProducer.removeMetrics(queueName);
return observers.remove(queueName) != null;
return observers.remove(queueName) != null || clientObservers.remove(queueName) != null;
}

@Override
Expand All @@ -106,8 +113,9 @@ public void shutdown() {
private void updateWorkerQueueState() {
if (running) {
try {
final Map<String, RabbitMQQueueStatus> queueStates = new HashMap<>(monitor.getWorkerQueueStates());
final Map<String, RabbitMQQueueStatus> queueStates = new HashMap<>(monitor.getQueueStates());
observers.forEach((q, v) -> updateWorkerQueueState(q, queueStates.get(q)));
clientObservers.forEach((q, v) -> updateClientQueueState(v, queueStates));
} catch (final RuntimeException e) {
LOG.error("Runtime error during updateWorkerQueueState", e);
}
Expand All @@ -120,6 +128,10 @@ private void updateWorkerQueueState(final String queueName, final RabbitMQQueueS
});
}

private void updateClientQueueState(final ClientQueueObserver observer, final Map<String, RabbitMQQueueStatus> queueStates) {
queueStates.entrySet().stream().filter(e -> observer.filter(e.getKey())).forEach(e -> observer.onClientQueueUpdate(e.getKey(), e.getValue()));
}

private static class WorkerSizeObserverComposite implements WorkerSizeObserver {
private final List<WorkerSizeObserver> observers = new ArrayList<>();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@
*/
class PriorityTaskSchedulerMetrics {

private static final String METRIC_PREFIX = "aer.taskmanager.client.queue";
private static final String METRIC_NAME = "aer.taskmanager.queue.usage";
private static final String DESCRIPTION = "Number of tasks running on client queues";

private final Map<String, ObservableDoubleGauge> usageMetrics = new HashMap<>();
Expand Down Expand Up @@ -64,7 +64,7 @@ private ObservableDoubleGauge createMetric(final IntSupplier countSupplier, fina
final Attributes queueAttributes = OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", state);

return OpenTelemetryMetrics.METER
.gaugeBuilder(METRIC_PREFIX)
.gaugeBuilder(METRIC_NAME)
.setDescription(DESCRIPTION)
.buildWithCallback(result -> result.record(countSupplier.getAsInt(), queueAttributes));
}
Expand Down
Loading
Loading