This is an automated email from the ASF dual-hosted git repository.
cschneider pushed a commit to branch master
in repository
https://gitbox.apache.org/repos/asf/sling-org-apache-sling-distribution-journal.git
The following commit(s) were added to refs/heads/master by this push:
new 27105b7 SLING-13135 - Compute queue size in background (#184)
27105b7 is described below
commit 27105b7d5810f766dd09db04d0d255bc6b62500d
Author: Christian Schneider <[email protected]>
AuthorDate: Thu Mar 12 14:18:20 2026 +0100
SLING-13135 - Compute queue size in background (#184)
* SLING-13135 - Compute queue size in background
* SLING-13135 - Fixes from review
* SLING-13135 - Add metric for queue size computation
* SLING-13135 - Fixes from review
* SLING-13135 - Fixes from review
* SLING-13135 - Fixes from review
* SLING-13135 - Fixes from review
---
docs/metrics_overview.md | 9 +-
.../journal/impl/publisher/PublishMetrics.java | 2 +
.../queue/impl/PubQueueProviderFactoryImpl.java | 28 ++++--
.../journal/queue/impl/PubQueueProviderImpl.java | 61 ++++++++++++
.../journal/queue/impl/PubQueueProviderTest.java | 110 +++++++++++++++++++++
5 files changed, 199 insertions(+), 11 deletions(-)
diff --git a/docs/metrics_overview.md b/docs/metrics_overview.md
index 395e7f2..c9f2f6c 100644
--- a/docs/metrics_overview.md
+++ b/docs/metrics_overview.md
@@ -45,7 +45,14 @@ 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)
+- **Description**: Current size of the queue (maximum queue size). 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.
+
+#### `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.
- **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/PublishMetrics.java
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java
index 5692450..d6d1bec 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
@@ -44,6 +44,8 @@ public class PublishMetrics {
private static final String QUEUE_ACCESS_ERROR_COUNT = PUB_COMPONENT +
"queue_access_error_count";
private static final String SUBSCRIBER_COUNT = PUB_COMPONENT +
"subscriber_count";
private static final String QUEUE_SIZE = PUB_COMPONENT + "queue_size";
+ /** 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 List<Tag> tags;
private final MetricsService metricsService;
diff --git
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
index 671cf77..2c4f001 100644
---
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
+++
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
@@ -18,32 +18,40 @@
*/
package org.apache.sling.distribution.journal.queue.impl;
+import org.apache.sling.commons.metrics.MetricsService;
import org.apache.sling.distribution.journal.queue.CacheCallback;
import org.apache.sling.distribution.journal.queue.PubQueueProvider;
import org.apache.sling.distribution.journal.queue.PubQueueProviderFactory;
import org.osgi.framework.BundleContext;
+import org.osgi.service.component.annotations.Activate;
import org.osgi.service.component.annotations.Component;
import org.osgi.service.component.annotations.Reference;
+import org.osgi.service.component.annotations.ReferenceCardinality;
import org.osgi.service.event.EventAdmin;
@Component
public class PubQueueProviderFactoryImpl implements PubQueueProviderFactory {
-
- @Reference
- private EventAdmin eventAdmin;
- @Reference
- private QueueErrors queueErrors;
-
- private BundleContext context;
+ private final EventAdmin eventAdmin;
+ private final QueueErrors queueErrors;
+ private final MetricsService metricsService;
+ private final BundleContext context;
- public void activate(BundleContext context) {
+ @Activate
+ public PubQueueProviderFactoryImpl(
+ @Reference EventAdmin eventAdmin,
+ @Reference QueueErrors queueErrors,
+ @Reference(cardinality = ReferenceCardinality.OPTIONAL)
MetricsService metricsService,
+ BundleContext context) {
+ this.eventAdmin = eventAdmin;
+ this.queueErrors = queueErrors;
+ this.metricsService = metricsService;
this.context = context;
}
@Override
public PubQueueProvider create(CacheCallback callback) {
- return new PubQueueProviderImpl(eventAdmin, queueErrors, callback,
context);
+ return new PubQueueProviderImpl(eventAdmin, queueErrors, callback,
context, metricsService);
}
-
+
}
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 3b2d74c..181efc8 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
@@ -20,6 +20,7 @@ package org.apache.sling.distribution.journal.queue.impl;
import static
org.apache.sling.commons.scheduler.Scheduler.PROPERTY_SCHEDULER_CONCURRENT;
import static
org.apache.sling.commons.scheduler.Scheduler.PROPERTY_SCHEDULER_PERIOD;
+import static
org.apache.sling.distribution.journal.metrics.TaggedMetrics.getMetricName;
import java.util.Dictionary;
import java.util.HashSet;
@@ -28,6 +29,9 @@ import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
@@ -35,12 +39,16 @@ import javax.annotation.ParametersAreNonnullByDefault;
import org.apache.commons.io.IOUtils;
import org.apache.commons.lang3.StringUtils;
+import org.apache.sling.commons.metrics.MetricsService;
+import org.apache.sling.commons.metrics.Timer;
import org.apache.sling.commons.scheduler.Scheduler;
import org.apache.sling.distribution.journal.MessageInfo;
import org.apache.sling.distribution.journal.impl.discovery.State;
import
org.apache.sling.distribution.journal.impl.publisher.PackageQueuedNotifier;
import org.apache.sling.distribution.journal.messages.PackageStatusMessage;
import
org.apache.sling.distribution.journal.messages.PackageStatusMessage.Status;
+import org.apache.sling.distribution.journal.impl.publisher.PublishMetrics;
+import org.apache.sling.distribution.journal.metrics.Tag;
import org.apache.sling.distribution.journal.queue.CacheCallback;
import org.apache.sling.distribution.journal.queue.OffsetQueue;
import org.apache.sling.distribution.journal.queue.PubQueueProvider;
@@ -62,6 +70,8 @@ public class PubQueueProviderImpl implements
PubQueueProvider, Runnable {
*/
private static final int CLEANUP_THRESHOLD = 10_000;
+ static final long QUEUE_SIZE_REFRESH_INTERVAL_SECONDS = 30;
+
private static final Logger LOG =
LoggerFactory.getLogger(PubQueueProviderImpl.class);
private final PackageQueuedNotifier queuedNotifier;
@@ -77,13 +87,33 @@ 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 ScheduledExecutorService queueSizeExecutor;
+
+ @Nullable
+ private final MetricsService metricsService;
+
private ServiceRegistration<?> reg;
public PubQueueProviderImpl(EventAdmin eventAdmin, QueueErrors
queueErrors, CacheCallback callback, BundleContext context) {
+ this(eventAdmin, queueErrors, callback, context, null);
+ }
+
+ public PubQueueProviderImpl(EventAdmin eventAdmin, QueueErrors
queueErrors, CacheCallback callback, BundleContext context,
+ @Nullable MetricsService metricsService) {
queuedNotifier = new PackageQueuedNotifier(eventAdmin);
this.queueErrors = queueErrors;
this.callback = callback;
+ this.metricsService = metricsService;
cache = newCache();
+ queueSizeExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
+ Thread t = new Thread(r, "queue-size-refresh");
+ t.setDaemon(true);
+ return t;
+ });
+ queueSizeExecutor.scheduleAtFixedRate(this::refreshQueueSizes,
+ 0, QUEUE_SIZE_REFRESH_INTERVAL_SECONDS, TimeUnit.SECONDS);
startCleanupTask(context);
LOG.info("Started Publisher queue provider service");
}
@@ -99,6 +129,7 @@ public class PubQueueProviderImpl implements
PubQueueProvider, Runnable {
@Override
public void close() {
+ queueSizeExecutor.shutdownNow();
PubQueueCache queueCache = this.cache;
if (queueCache != null) {
queueCache.close();
@@ -209,6 +240,10 @@ public class PubQueueProviderImpl implements
PubQueueProvider, Runnable {
@Override
public int getMaxQueueSize(String pubAgentName) {
+ return cachedQueueSizes.computeIfAbsent(pubAgentName, key -> 0);
+ }
+
+ private int computeMaxQueueSize(String pubAgentName) {
Optional<Long> minOffset = getMinEditableQueueOffset(pubAgentName);
if (minOffset.isPresent()) {
return getOffsetQueue(pubAgentName,
minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize();
@@ -217,6 +252,32 @@ public class PubQueueProviderImpl implements
PubQueueProvider, Runnable {
}
}
+ /**
+ * Package-private for tests. Triggers an immediate synchronous refresh of
cached queue sizes.
+ */
+ void triggerQueueSizeRefreshForTest() {
+ refreshQueueSizes();
+ }
+
+ private void refreshQueueSizes() {
+ for (String agentName : cachedQueueSizes.keySet()) {
+ try {
+ long startNanos = System.nanoTime();
+ int size = computeMaxQueueSize(agentName);
+ long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000;
+ cachedQueueSizes.put(agentName, size);
+
+ if (metricsService != null) {
+ Timer timer = metricsService.timer(getMetricName(
+ PublishMetrics.QUEUE_SIZE_COMPUTATION_DURATION,
Tag.of("pub_name", agentName)));
+ timer.update(elapsedMs, TimeUnit.MILLISECONDS);
+ }
+ } catch (Exception e) {
+ LOG.warn("Failed to refresh queue size for agent {}",
agentName, e);
+ }
+ }
+ }
+
private Optional<Long> getMinEditableQueueOffset(String pubAgentName) {
return callback.getSubscribedAgentIds(pubAgentName).stream()
.filter(subAgentName -> isEditable(pubAgentName, subAgentName))
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 008610e..dec721c 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
@@ -20,6 +20,8 @@ package org.apache.sling.distribution.journal.queue.impl;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.equalTo;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.atLeast;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
@@ -29,6 +31,7 @@ import java.io.Closeable;
import java.io.IOException;
import java.lang.management.ManagementFactory;
import java.util.Collections;
+import java.util.concurrent.TimeUnit;
import java.util.Iterator;
import java.util.List;
import java.util.Set;
@@ -42,6 +45,8 @@ import javax.management.ObjectInstance;
import javax.management.ObjectName;
import javax.management.ReflectionException;
+import org.apache.sling.commons.metrics.MetricsService;
+import org.apache.sling.commons.metrics.Timer;
import org.apache.sling.distribution.journal.HandlerAdapter;
import org.apache.sling.distribution.journal.MessageHandler;
import org.apache.sling.distribution.journal.MessageInfo;
@@ -176,6 +181,8 @@ public class PubQueueProviderTest {
State state = Mockito.mock(State.class);
when(state.isEditable()).thenReturn(true);
when(callback.getState(Mockito.eq(PUB1_AGENT_NAME),
Mockito.anyString())).thenReturn(state);
+ queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+ queueProvider.triggerQueueSizeRefreshForTest();
int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
assertThat(size, equalTo(2));
}
@@ -190,6 +197,27 @@ public class PubQueueProviderTest {
int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
assertThat(size, equalTo(0));
}
+
+ @Test
+ public void testQueueSizeIsCached() {
+ handler.handle(info(0L), packageMessage("packageid1",
PUB1_AGENT_NAME));
+ handler.handle(info(1L), packageMessage("packageid2",
PUB2_AGENT_NAME));
+ handler.handle(info(2L), packageMessage("packageid3",
PUB1_AGENT_NAME));
+
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, null));
+ State state = Mockito.mock(State.class);
+ when(state.isEditable()).thenReturn(true);
+ when(callback.getState(Mockito.eq(PUB1_AGENT_NAME),
Mockito.anyString())).thenReturn(state);
+
+ queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+ queueProvider.triggerQueueSizeRefreshForTest();
+ assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+
+ handler.handle(info(3L), packageMessage("packageid4",
PUB1_AGENT_NAME));
+
+ assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+ }
@SuppressWarnings("null")
@Test
@@ -241,6 +269,88 @@ public class PubQueueProviderTest {
assertThat(queueSize(), equalTo(1));
}
+ @Test
+ public void testQueueSizeNewAgentReturnsZeroUntilRefresh() {
+ handler.handle(info(0L), packageMessage("packageid1",
PUB1_AGENT_NAME));
+ handler.handle(info(1L), packageMessage("packageid2",
PUB2_AGENT_NAME));
+ handler.handle(info(2L), packageMessage("packageid3",
PUB1_AGENT_NAME));
+
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, null));
+ State state = Mockito.mock(State.class);
+ when(state.isEditable()).thenReturn(true);
+ when(callback.getState(Mockito.eq(PUB1_AGENT_NAME),
Mockito.anyString())).thenReturn(state);
+
+ assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0));
+ queueProvider.triggerQueueSizeRefreshForTest();
+ assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+ }
+
+ @Test
+ public void testQueueSizeWithMetricsServiceRecordsTimer() {
+ MetricsService mockMetrics = mock(MetricsService.class);
+ Timer mockTimer = mock(Timer.class);
+ when(mockMetrics.timer(Mockito.anyString())).thenReturn(mockTimer);
+
+ QueueErrors queueErrors = mock(QueueErrors.class);
+ PubQueueProviderImpl providerWithMetrics = new PubQueueProviderImpl(
+ eventAdmin, queueErrors, callback, context, mockMetrics);
+ MessageHandler<PackageMessage> handler2 =
handlerCaptor.getAllValues().get(1);
+
+ handler2.handle(info(0L), packageMessage("packageid1",
PUB1_AGENT_NAME));
+ handler2.handle(info(1L), packageMessage("packageid2",
PUB2_AGENT_NAME));
+ handler2.handle(info(2L), packageMessage("packageid3",
PUB1_AGENT_NAME));
+
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, null));
+ State state = Mockito.mock(State.class);
+ when(state.isEditable()).thenReturn(true);
+ when(callback.getState(Mockito.eq(PUB1_AGENT_NAME),
Mockito.anyString())).thenReturn(state);
+
+ providerWithMetrics.getMaxQueueSize(PUB1_AGENT_NAME);
+ providerWithMetrics.triggerQueueSizeRefreshForTest();
+ verify(mockTimer).update(anyLong(), eq(TimeUnit.MILLISECONDS));
+ providerWithMetrics.close();
+ }
+
+ @Test
+ public void testRefreshQueueSizesHandlesExceptionForOneAgent() {
+ handler.handle(info(0L), packageMessage("packageid1",
PUB1_AGENT_NAME));
+ handler.handle(info(1L), packageMessage("packageid2",
PUB2_AGENT_NAME));
+ handler.handle(info(2L), packageMessage("packageid3",
PUB1_AGENT_NAME));
+ when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME))
+ .thenThrow(new RuntimeException("test"));
+ when(callback.getSubscribedAgentIds(PUB2_AGENT_NAME))
+ .thenReturn(Collections.singleton("sub1"));
+ when(callback.getQueueState(Mockito.eq(PUB2_AGENT_NAME),
Mockito.eq("sub1")))
+ .thenReturn(new QueueState(0, -1, 0, null));
+ State state = Mockito.mock(State.class);
+ when(state.isEditable()).thenReturn(true);
+ when(callback.getState(Mockito.eq(PUB2_AGENT_NAME),
Mockito.eq(SUB_AGENT_NAME))).thenReturn(state);
+
+ queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+ queueProvider.getMaxQueueSize(PUB2_AGENT_NAME);
+ queueProvider.triggerQueueSizeRefreshForTest();
+ assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME), equalTo(1));
+ }
+
+ @Test
+ public void testQueueSizeWithNoEditableSubscribers() {
+ handler.handle(info(0L), packageMessage("packageid1",
PUB1_AGENT_NAME));
+ handler.handle(info(1L), packageMessage("packageid2",
PUB2_AGENT_NAME));
+ handler.handle(info(2L), packageMessage("packageid3",
PUB1_AGENT_NAME));
+
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, null));
+ State state = Mockito.mock(State.class);
+ when(state.isEditable()).thenReturn(false);
+ when(callback.getState(Mockito.eq(PUB1_AGENT_NAME),
Mockito.anyString())).thenReturn(state);
+
+ queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+ queueProvider.triggerQueueSizeRefreshForTest();
+ assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0));
+ }
+
private int queueSize() {
return queueProvider.getOffsetQueue(PUB1_AGENT_NAME, 0).getSize();
}