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
11 changes: 0 additions & 11 deletions doc/telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -114,17 +114,6 @@ Third the `aer.rabbitmq.worker.*` metrics are the values as received from the R
These metrics could also be obtained by directly reading the RabbitMQ api, but specifically for the usage don't require additional logic to get the usage metrics.
In general these metrics should report the same values, but due to timing (e.g. the moment the measure is taken) there can be differences.

##### Deprecated metrics

The following metrics have been replaced by the more standardized naming mentioned above

| Metric name | type | description | Replaced by |
|-------------------------------------------------------|-----------|----------------------------------------------------------------------|--------------------------------------------|
| `aer.taskmanager.worker_size`<sup>1</sup> | gauge | The sum of idle workers + occupied workers. | `aer.taskmanager.workerpool.worker.limit` |
| `aer.taskmanager.current_worker_size`<sup>1</sup> | gauge | The number of workers based on what RabbitMQ reports. | `aer.rabbitmq.worker.limit` |
| `aer.taskmanager.running_worker_size`<sup>1</sup> | gauge | The number of workers that are occupied. | `aer.taskmanager.workerpool.worker..usage` |
| `aer.taskmanager.running_client_size`<sup>3</sup> | gauge | The number of workers that are occupied for a specific client queue. | `aer.taskmanager.client.queue.usage` |

##### Metric attributes

The workers have different attributes to distinguish specific metrics.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,8 @@ public static long longValue(final Map<String, Object> messageMetaData, final St
}

public static String stringValue(final Map<String, Object> messageMetaData, final String key) {
return Optional.ofNullable(messageMetaData.get(key))
return Optional.ofNullable(messageMetaData)
.map(m -> m.get(key))
.filter(LongString.class::isInstance)
.map(t -> new String(((LongString) t).getBytes()))
.orElse("");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,17 +17,18 @@
package nl.aerius.taskmanager;

import java.util.concurrent.Semaphore;
import java.util.concurrent.locks.ReentrantLock;

import nl.aerius.taskmanager.adaptor.WorkerSizeObserver;

/**
* Class to be used at startup. The Scheduler should not start before it is known how many messages are still on the queue.
* This to register any work that is still on the queue and to properly calculate load metrics.
* Because the Task Manager is not aware of the tasks already on the queue and therefore otherwise these messages won't be counted in the metrics.
* This can result in the metrics being skewed, and thereby negatively reporting load metrics.
* Class to be used at startup. The Scheduler should not start before the number of messages on the queue is zero.
* Because the Task Manager has no information of the tasks already on the queue and therefore there is no tracking information of those messages.
* As all tracking information only lives in memory and is reset when the Task Manager is restarted.
*/
public class StartupGuard implements WorkerSizeObserver {

private final ReentrantLock lock = new ReentrantLock();
private final Semaphore openSemaphore = new Semaphore(0);

private boolean open;
Expand All @@ -48,11 +49,14 @@ public void waitForOpen() throws InterruptedException {

@Override
public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) {
synchronized (openSemaphore) {
lock.lock();
try {
if (!open) {
open = true;
openSemaphore.release();
}
} finally {
lock.unlock();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -174,8 +174,6 @@ public TaskScheduleBucket(final QueueConfig queueConfig) throws InterruptedExcep
workerSizeObserverProxy.addObserver(workerQueueName, wzo);
}
workerProducer.start();
// Set up metrics
WorkerPoolMetrics.setupMetrics(workerPool, workerQueueName);

dispatcher = new TaskDispatcher(workerQueueName, taskScheduler, workerPool);
executorService.execute(() -> {
Expand Down Expand Up @@ -234,7 +232,6 @@ public void addTaskConsumerIfAbsent(final QueueConfig queueConfig) {
});
}


/**
* Removes a task consumer with the given queue name.
*
Expand All @@ -250,7 +247,6 @@ public void shutdown() {
dispatcher.shutdown();
workerProducer.shutdown();
taskManagerMetrics.remove(workerQueueName);
WorkerPoolMetrics.removeMetrics(workerQueueName);
taskConsumers.forEach((k, v) -> v.shutdown());
}

Expand Down

This file was deleted.

Loading