This is an automated email from the ASF dual-hosted git repository. cschneider pushed a commit to branch SLING-13136-2 in repository https://gitbox.apache.org/repos/asf/sling-org-apache-sling-distribution-journal.git
commit 0ab8f7fef97aea679dbf5a180652551914596071 Author: Christian Schneider <[email protected]> AuthorDate: Mon Mar 23 17:44:06 2026 +0100 SLING-13136 - Provide separate metric for largest (non) clearable queue --- README.md | 2 +- docs/metrics_overview.md | 13 ++++--- .../impl/publisher/DistributionPublisher.java | 5 ++- .../journal/impl/publisher/PublishMetrics.java | 15 +++++++- .../journal/queue/PubQueueProvider.java | 9 +++-- .../journal/queue/impl/PubQueueProviderImpl.java | 45 +++++++++++++++------- .../impl/publisher/DistributionPublisherTest.java | 5 ++- .../journal/queue/impl/PubQueueProviderTest.java | 30 ++++++++------- 8 files changed, 81 insertions(+), 43 deletions(-) diff --git a/README.md b/README.md index 3e72ccc..8e6044f 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,7 @@ Publisher metrics (prefixed with `sling_distribution_journal_publisher_`) track - **Package Export**: `sling_distribution_journal_publisher_exported_package_size` (histogram) - **Request Handling**: `sling_distribution_journal_publisher_accepted_requests`, `sling_distribution_journal_publisher_dropped_requests` (meters) - **Package Building**: `sling_distribution_journal_publisher_build_package_duration`, `sling_distribution_journal_publisher_enqueue_package_duration` (timers) -- **Queue Operations**: `sling_distribution_journal_publisher_queue_size` (gauge), `sling_distribution_journal_publisher_queue_cache_fetch_count`, `sling_distribution_journal_publisher_queue_access_error_count` (counters) +- **Queue Operations**: `sling_distribution_journal_publisher_queue_size` (gauge, tagged `pub_name` and `clearable`), `sling_distribution_journal_publisher_queue_cache_fetch_count`, `sling_distribution_journal_publisher_queue_access_error_count` (counters) - **Subscriber Discovery**: `sling_distribution_journal_publisher_subscriber_count` (gauge) ### Subscriber Metrics diff --git a/docs/metrics_overview.md b/docs/metrics_overview.md index 2935746..ebed07f 100644 --- a/docs/metrics_overview.md +++ b/docs/metrics_overview.md @@ -4,9 +4,11 @@ This document provides a comprehensive overview of all metrics in the Apache Sli ## Publisher Metrics -All publisher metrics are prefixed with `sling_distribution_journal_publisher_` and include the following tag: +Most publisher metrics are prefixed with `sling_distribution_journal_publisher_` and include: - `pub_name`: Name of the publish agent +The `queue_size` gauge additionally includes `clearable` (see below). + ### Package Export Metrics #### `sling_distribution_journal_publisher_exported_package_size` (Histogram) @@ -45,14 +47,15 @@ All publisher metrics are prefixed with `sling_distribution_journal_publisher_` #### `sling_distribution_journal_publisher_queue_size` (Gauge) - **Type**: Gauge -- **Description**: Current size of the queue (maximum queue size). The backlog is measured from the minimum `lastProcessedOffset` among **clearable** subscribers (those with a non-null clear callback on `QueueState`). Values are cached and refreshed in the background every 30 seconds. -- **Tags**: `pub_name` -- **Staleness**: The value can be up to ~30 seconds stale under normal conditions. When `computeMaxQueueSize()` scales linearly (O(n)) with queue size and takes longer than the refresh interval, staleness increases: refresh cycles can back up, and the displayed value may lag further behind the true queue size. See `sling_distribution_journal_publisher_queue_size_computation_duration` for monitoring computation cost. +- **Description**: Max backlog depth for a subscriber cohort, measured from the minimum `lastProcessedOffset` in that cohort (tail size in the journal cache). **Clearable** cohort: subscribers with a non-null clear callback on `QueueState`. **Non-clearable** cohort: subscribers with a `QueueState` and a null clear callback. Two time series per publisher (see tags). Values are cached and refreshed in the background every 30 seconds. Publisher throttling uses the **clearable** series only [...] +- **Tags**: `pub_name`, `clearable` (`true` or `false`) +- **Migration**: Previously this metric used only `pub_name`. Series are now distinguished by `clearable`; update dashboards and alerts to include the `clearable` tag (e.g. `clearable=true` for the former single series). +- **Staleness**: The value can be up to ~30 seconds stale under normal conditions. When queue size computation scales linearly (O(n)) with queue size and takes longer than the refresh interval, staleness increases: refresh cycles can back up, and the displayed value may lag further behind the true queue size. See `sling_distribution_journal_publisher_queue_size_computation_duration` for monitoring computation cost. #### `sling_distribution_journal_publisher_queue_size_computation_duration` (Timer) - **Type**: Timer - **Unit**: Milliseconds -- **Description**: Duration of computing the max queue size per agent during background refresh. The computation is O(n) with queue size; slow values (>30s) indicate increased staleness. +- **Description**: Duration of computing **both** clearable and non-clearable max queue sizes for an agent during each background refresh. The computation is O(n) with queue size; slow values (>30s) indicate increased staleness. - **Tags**: `pub_name` #### `sling_distribution_journal_publisher_queue_cache_fetch_count` (Counter) diff --git a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java index 5829126..33ba863 100644 --- a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java +++ b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java @@ -140,7 +140,8 @@ public class DistributionPublisher implements DistributionAgent { requireNonNull(metricsService); this.publishMetrics = new PublishMetrics(metricsService, pubAgentName); this.pubQueueProvider = pubQueueProvider; - this.publishMetrics.queueSize(() -> pubQueueProvider.getMaxQueueSize(pubAgentName)); + this.publishMetrics.queueSizeByClearable(() -> pubQueueProvider.getMaxQueueSize(pubAgentName, true), true); + this.publishMetrics.queueSizeByClearable(() -> pubQueueProvider.getMaxQueueSize(pubAgentName, false), false); distLog = new DefaultDistributionLog(pubAgentName, this.getClass(), DefaultDistributionLog.LogLevel.INFO); distributionLogEventListener = new DistributionLogEventListener(context, distLog, pubAgentName); @@ -212,7 +213,7 @@ public class DistributionPublisher implements DistributionAgent { distLog.info(msg); return new SimpleDistributionResponse(DistributionRequestState.DROPPED, msg); } - int queueSize = pubQueueProvider.getMaxQueueSize(pubAgentName); + int queueSize = pubQueueProvider.getMaxQueueSize(pubAgentName, true); int sleepMs = getSleepTime(queueSize); sleep(sleepMs); final PackageMessage pkg = buildPackage(resourceResolver, request); diff --git a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java index d6d1bec..628c631 100644 --- a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java +++ b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java @@ -33,6 +33,7 @@ import org.apache.sling.distribution.journal.metrics.Tag; public class PublishMetrics { private static final String TAG_AGENT_NAME = "pub_name"; + private static final String TAG_CLEARABLE = "clearable"; public static final String PUB_COMPONENT = "distribution.journal.publisher."; private static final String EXPORTED_PACKAGE_SIZE = PUB_COMPONENT + "exported_package_size"; @@ -47,10 +48,12 @@ public class PublishMetrics { /** Metric name for queue size computation duration (use with Tag.of("pub_name", agentName)). */ public static final String QUEUE_SIZE_COMPUTATION_DURATION = PUB_COMPONENT + "queue_size_computation_duration"; + private final String pubAgentName; private final List<Tag> tags; private final MetricsService metricsService; public PublishMetrics(MetricsService metricsService, String pubAgentName) { + this.pubAgentName = pubAgentName; this.tags = Arrays.asList(Tag.of(TAG_AGENT_NAME, pubAgentName)); this.metricsService = metricsService; } @@ -122,8 +125,16 @@ public class PublishMetrics { metricsService.gauge(getMetricName(SUBSCRIBER_COUNT, tags), subscriberCountCallback); } - public void queueSize(Supplier<Integer> queueSizeCallback) { - metricsService.gauge(getMetricName(QUEUE_SIZE, tags), queueSizeCallback); + /** + * Gauge of max queue backlog for subscribers in the given cohort. + * + * @param clearable {@code true} for clearable subscribers, {@code false} for non-clearable + */ + public void queueSizeByClearable(Supplier<Integer> queueSizeCallback, boolean clearable) { + List<Tag> queueSizeTags = Arrays.asList( + Tag.of(TAG_AGENT_NAME, pubAgentName), + Tag.of(TAG_CLEARABLE, Boolean.toString(clearable))); + metricsService.gauge(getMetricName(QUEUE_SIZE, queueSizeTags), queueSizeCallback); } } diff --git a/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java b/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java index 3da54df..7cefb90 100644 --- a/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java +++ b/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java @@ -38,11 +38,14 @@ public interface PubQueueProvider extends Closeable { DistributionQueue getQueue(String pubAgentName, String queueName); /** - * Get maximum size of all queues for a pubAgentName + * Maximum backlog depth for a subscriber cohort (min {@code lastProcessedOffset} in cohort, then journal tail size). + * * @param pubAgentName name of the pub agent - * @return max size of all queues or 0 if there are none + * @param clearable {@code true} for clearable subscribers (non-null clear callback); {@code false} for non-clearable + * (queue state present, null clear callback). Publisher throttling uses {@code clearable == true}. + * @return max size for that cohort or 0 if there are none */ - int getMaxQueueSize(String pubAgentName); + int getMaxQueueSize(String pubAgentName, boolean clearable); @Nonnull OffsetQueue<DistributionQueueItem> getOffsetQueue(String pubAgentName, long minOffset); diff --git a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java index 2b47f57..8cb8c33 100644 --- a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java +++ b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java @@ -86,7 +86,7 @@ public class PubQueueProviderImpl implements PubQueueProvider, Runnable { */ private final Map<String, OffsetQueue<Long>> errorQueues = new ConcurrentHashMap<>(); - private final Map<String, Integer> cachedQueueSizes = new ConcurrentHashMap<>(); + private final Map<String, CachedMaxQueueSizes> cachedMaxQueueSizes = new ConcurrentHashMap<>(); private final ScheduledExecutorService queueSizeExecutor; @@ -238,17 +238,17 @@ public class PubQueueProviderImpl implements PubQueueProvider, Runnable { @Override - public int getMaxQueueSize(String pubAgentName) { - return cachedQueueSizes.computeIfAbsent(pubAgentName, key -> 0); + public int getMaxQueueSize(String pubAgentName, boolean clearable) { + CachedMaxQueueSizes entry = cachedMaxQueueSizes.computeIfAbsent(pubAgentName, k -> new CachedMaxQueueSizes()); + return clearable ? entry.clearable : entry.nonClearable; } - private int computeMaxQueueSize(String pubAgentName) { - Optional<Long> minOffset = getMinClearableQueueOffset(pubAgentName); + private int computeMaxQueueSize(String pubAgentName, boolean clearableCohort) { + Optional<Long> minOffset = getMinQueueOffsetForCohort(pubAgentName, clearableCohort); if (minOffset.isPresent()) { return getOffsetQueue(pubAgentName, minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize(); - } else { - return 0; } + return 0; } /** @@ -259,12 +259,17 @@ public class PubQueueProviderImpl implements PubQueueProvider, Runnable { } private void refreshQueueSizes() { - for (String agentName : cachedQueueSizes.keySet()) { + for (String agentName : cachedMaxQueueSizes.keySet()) { try { long startNanos = System.nanoTime(); - int size = computeMaxQueueSize(agentName); + int clearableSize = computeMaxQueueSize(agentName, true); + int nonClearableSize = computeMaxQueueSize(agentName, false); long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000; - cachedQueueSizes.put(agentName, size); + CachedMaxQueueSizes entry = cachedMaxQueueSizes.get(agentName); + if (entry != null) { + entry.clearable = clearableSize; + entry.nonClearable = nonClearableSize; + } if (metricsService != null) { Timer timer = metricsService.timer(getMetricName( @@ -277,16 +282,23 @@ public class PubQueueProviderImpl implements PubQueueProvider, Runnable { } } - private Optional<Long> getMinClearableQueueOffset(String pubAgentName) { + private Optional<Long> getMinQueueOffsetForCohort(String pubAgentName, boolean clearableCohort) { return callback.getSubscribedAgentIds(pubAgentName).stream() - .filter(subAgentName -> isClearable(pubAgentName, subAgentName)) + .filter(subAgentName -> inClearableCohort(pubAgentName, subAgentName, clearableCohort)) .map(subAgentName -> lastProcessedOffset(pubAgentName, subAgentName)) .min(Long::compare); } - private boolean isClearable(String pubAgentName, String subAgentName) { + /** + * @param clearableCohort {@code true} for clearable subscribers, {@code false} for non-clearable (state present, no clear callback). + */ + private boolean inClearableCohort(String pubAgentName, String subAgentName, boolean clearableCohort) { QueueState queueState = callback.getQueueState(pubAgentName, subAgentName); - return queueState != null && queueState.getClearCallback() != null; + if (queueState == null) { + return false; + } + boolean clearable = queueState.getClearCallback() != null; + return clearableCohort ? clearable : !clearable; } private long lastProcessedOffset(String pubAgentName, String subAgentName) { @@ -309,4 +321,9 @@ public class PubQueueProviderImpl implements PubQueueProvider, Runnable { return new PubQueueCache(queuedNotifier, callback); } + private static final class CachedMaxQueueSizes { + volatile int clearable; + volatile int nonClearable; + } + } diff --git a/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java b/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java index 200d564..e65f4c4 100644 --- a/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java +++ b/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java @@ -33,6 +33,7 @@ import static org.junit.Assert.assertNull; import static org.junit.Assert.fail; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -175,7 +176,7 @@ public class DistributionPublisherTest { @Test public void testQueueSizeLimitHalf() throws IOException, DistributionException { int queueSize = DEFAULT_QUEUE_SIZE_LIMIT + DEFAULT_QUEUE_SIZE_LIMIT / 2; - when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1)).thenReturn(queueSize); + when(pubQueueProvider.getMaxQueueSize(eq(PUB1AGENT1), eq(true))).thenReturn(queueSize); DistributionRequest request = new SimpleDistributionRequest(DistributionRequestType.ADD, "/test"); long time = distribute(request); assertThat(time, greaterThanOrEqualTo(500L)); @@ -185,7 +186,7 @@ public class DistributionPublisherTest { @Test public void testDoubleQueueSizeLimitReached() throws IOException, DistributionException { int queueSize = DEFAULT_QUEUE_SIZE_LIMIT * 2; - when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1)).thenReturn(queueSize); + when(pubQueueProvider.getMaxQueueSize(eq(PUB1AGENT1), eq(true))).thenReturn(queueSize); DistributionRequest request = new SimpleDistributionRequest(DistributionRequestType.ADD, "/test"); long time = distribute(request); assertThat(time, greaterThanOrEqualTo(1000L)); diff --git a/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java b/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java index 2f25602..05108b0 100644 --- a/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java +++ b/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java @@ -178,10 +178,11 @@ public class PubQueueProviderTest { when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(Collections.singleton("sub1")); when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class))); - queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); + queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); queueProvider.triggerQueueSizeRefreshForTest(); - int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); + int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); assertThat(size, equalTo(2)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, false), equalTo(0)); } @Test @@ -191,7 +192,7 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, null)); - int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); + int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); assertThat(size, equalTo(0)); } @@ -204,13 +205,13 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class))); - queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); + queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); queueProvider.triggerQueueSizeRefreshForTest(); - assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), equalTo(2)); handler.handle(info(3L), packageMessage("packageid4", PUB1_AGENT_NAME)); - assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), equalTo(2)); } @SuppressWarnings("null") @@ -272,9 +273,9 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class))); - assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), equalTo(0)); queueProvider.triggerQueueSizeRefreshForTest(); - assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), equalTo(2)); } @Test @@ -295,7 +296,7 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class))); - providerWithMetrics.getMaxQueueSize(PUB1_AGENT_NAME); + providerWithMetrics.getMaxQueueSize(PUB1_AGENT_NAME, true); providerWithMetrics.triggerQueueSizeRefreshForTest(); verify(mockTimer).update(anyLong(), eq(TimeUnit.MILLISECONDS)); providerWithMetrics.close(); @@ -313,10 +314,10 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB2_AGENT_NAME), Mockito.eq("sub1"))) .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class))); - queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); - queueProvider.getMaxQueueSize(PUB2_AGENT_NAME); + queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); + queueProvider.getMaxQueueSize(PUB2_AGENT_NAME, true); queueProvider.triggerQueueSizeRefreshForTest(); - assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME), equalTo(1)); + assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME, true), equalTo(1)); } @Test @@ -328,9 +329,10 @@ public class PubQueueProviderTest { when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), Mockito.any())) .thenReturn(new QueueState(0, -1, 0, null)); - queueProvider.getMaxQueueSize(PUB1_AGENT_NAME); + queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true); queueProvider.triggerQueueSizeRefreshForTest(); - assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), equalTo(0)); + assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, false), equalTo(2)); } private int queueSize() {
