diff --git a/doc/telemetry.md b/doc/telemetry.md index 70a1027..6c4e29c 100644 --- a/doc/telemetry.md +++ b/doc/telemetry.md @@ -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`1 | gauge | The sum of idle workers + occupied workers. | `aer.taskmanager.workerpool.worker.limit` | -| `aer.taskmanager.current_worker_size`1 | gauge | The number of workers based on what RabbitMQ reports. | `aer.rabbitmq.worker.limit` | -| `aer.taskmanager.running_worker_size`1 | gauge | The number of workers that are occupied. | `aer.taskmanager.workerpool.worker..usage` | -| `aer.taskmanager.running_client_size`3 | 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. diff --git a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/TaskMetrics.java b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/TaskMetrics.java index 0263d0e..f44e552 100644 --- a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/TaskMetrics.java +++ b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/TaskMetrics.java @@ -61,7 +61,8 @@ public static long longValue(final Map messageMetaData, final St } public static String stringValue(final Map 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(""); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java index e73aa7f..f8e547c 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java @@ -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; @@ -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(); } } } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java index 4258c42..55c6ed1 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java @@ -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(() -> { @@ -234,7 +232,6 @@ public void addTaskConsumerIfAbsent(final QueueConfig queueConfig) { }); } - /** * Removes a task consumer with the given queue name. * @@ -250,7 +247,6 @@ public void shutdown() { dispatcher.shutdown(); workerProducer.shutdown(); taskManagerMetrics.remove(workerQueueName); - WorkerPoolMetrics.removeMetrics(workerQueueName); taskConsumers.forEach((k, v) -> v.shutdown()); } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPoolMetrics.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPoolMetrics.java deleted file mode 100644 index 1e56627..0000000 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPoolMetrics.java +++ /dev/null @@ -1,98 +0,0 @@ -/* - * 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; - -import java.util.HashMap; -import java.util.Locale; -import java.util.Map; -import java.util.function.Function; - -import io.opentelemetry.api.common.Attributes; -import io.opentelemetry.api.metrics.ObservableDoubleGauge; - -import nl.aerius.taskmanager.metrics.OpenTelemetryMetrics; -import nl.aerius.taskmanager.metrics.UsageMetricsProvider; - -/** - * Set up metric collection for this worker pool with the given type name. - */ -public final class WorkerPoolMetrics { - - private static final Map REGISTERED_METRICS = new HashMap<>(); - - private enum WorkerPoolMetricType { - // @formatter:off - WORKER_SIZE(UsageMetricsProvider::getNumberOfWorkers, "Number of workers based on internal state of taskmanager"), - @Deprecated - CURRENT_WORKER_SIZE(WorkerPool::getReportedWorkerSize, - "Current number of workers according to taskmanager (deprecated replaced with 'aer.taskmanager.workerppol.worker.usage')"), - @Deprecated - RUNNING_WORKER_SIZE(UsageMetricsProvider::getNumberOfUsedWorkers, - "Used number of workers according to taskmanager (deprecated replaced with 'aer.taskmanager.workerppol.worker.usage')"); - // @formatter:on - - private final Function function; - private final String description; - - WorkerPoolMetricType(final Function function, final String description) { - this.function = function; - this.description = description; - } - - int getValue(final WorkerPool workerPool) { - return function.apply(workerPool); - } - - String getGaugeName() { - return "aer.taskmanager." + name().toLowerCase(Locale.ROOT); - } - - String getDescription() { - return description; - } - } - - private WorkerPoolMetrics() { - // Util-like class - } - - public static void setupMetrics(final WorkerPool workerPool, final String workerQueueName) { - final Attributes attributes = OpenTelemetryMetrics.workerAttributes(workerQueueName); - - for (final WorkerPoolMetricType metricType : WorkerPoolMetricType.values()) { - REGISTERED_METRICS.put(gaugeIdentifier(workerQueueName, metricType), - OpenTelemetryMetrics.METER.gaugeBuilder(metricType.getGaugeName()) - .setDescription(metricType.getDescription()) - .buildWithCallback( - result -> result.record(metricType.getValue(workerPool), attributes))); - } - } - - public static void removeMetrics(final String workerQueueName) { - for (final WorkerPoolMetricType metricType : WorkerPoolMetricType.values()) { - final String gaugeId = gaugeIdentifier(workerQueueName, metricType); - if (REGISTERED_METRICS.containsKey(gaugeId)) { - REGISTERED_METRICS.remove(gaugeId).close(); - } - } - } - - private static String gaugeIdentifier(final String workerQueueName, final WorkerPoolMetricType gaugeType) { - return workerQueueName + "_" + gaugeType.name(); - } - -}