This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 686bd4a40588 CAMEL-25109: camel-core - Scheduled poll consumer: fix
bugs found in a deep review (#27015)
686bd4a40588 is described below
commit 686bd4a405882a2c85595602090359629d5f2d89
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 22:44:28 2026 +0200
CAMEL-25109: camel-core - Scheduled poll consumer: fix bugs found in a deep
review (#27015)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../impl/ScheduledPollConsumerEdgeCasesTest.java | 226 +++++++++++++++++++++
.../DefaultScheduledPollConsumerScheduler.java | 5 +-
.../camel/support/ScheduledPollConsumer.java | 67 ++++--
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 9 +
4 files changed, 286 insertions(+), 21 deletions(-)
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
new file mode 100644
index 000000000000..102438842916
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
@@ -0,0 +1,226 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.impl;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.health.HealthCheck;
+import org.apache.camel.spi.HttpResponseAware;
+import org.apache.camel.support.DefaultScheduledPollConsumerScheduler;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class ScheduledPollConsumerEdgeCasesTest extends ContextTestSupport {
+
+ private static class MyConsumer extends MockScheduledPollConsumer {
+
+ private CountDownLatch pollStarted;
+ private CountDownLatch pollContinue;
+
+ MyConsumer(Endpoint endpoint, Exception exceptionToThrowOnPoll) {
+ super(endpoint, exceptionToThrowOnPoll);
+ }
+
+ @Override
+ protected int poll() throws Exception {
+ if (pollStarted != null) {
+ pollStarted.countDown();
+ pollContinue.await(10, TimeUnit.SECONDS);
+ }
+ return super.poll();
+ }
+
+ Throwable lastError() {
+ return getLastError();
+ }
+
+ Map<String, Object> lastErrorDetails() {
+ return getLastErrorDetails();
+ }
+ }
+
+ private static class MyHttpException extends Exception implements
HttpResponseAware {
+ MyHttpException() {
+ super("Service unavailable");
+ }
+
+ @Override
+ public int getHttpResponseCode() {
+ return 503;
+ }
+
+ @Override
+ public void setHttpResponseCode(int code) {
+ // noop
+ }
+
+ @Override
+ public String getHttpResponseStatus() {
+ return "Service Unavailable";
+ }
+
+ @Override
+ public void setHttpResponseStatus(String status) {
+ // noop
+ }
+ }
+
+ @Test
+ public void testChangedDelayUsedAfterRestart() {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"),
null);
+ consumer.setStartScheduler(false);
+ consumer.setDelay(20);
+ consumer.start();
+ consumer.stop();
+
+ consumer.setDelay(2000);
+ consumer.setInitialDelay(3000);
+ consumer.start();
+ try {
+ DefaultScheduledPollConsumerScheduler scheduler =
(DefaultScheduledPollConsumerScheduler) consumer.getScheduler();
+ assertEquals(2000, scheduler.getDelay());
+ assertEquals(3000, scheduler.getInitialDelay());
+ } finally {
+ consumer.stop();
+ }
+ }
+
+ @Test
+ public void testConcurrentUnscheduleTask() throws Exception {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"),
null);
+ List<Throwable> errors = new CopyOnWriteArrayList<>();
+
+ for (int i = 0; i < 30; i++) {
+ DefaultScheduledPollConsumerScheduler scheduler = new
DefaultScheduledPollConsumerScheduler();
+ scheduler.setCamelContext(context);
+ scheduler.setConcurrentConsumers(8);
+ scheduler.setPoolSize(8);
+ scheduler.setInitialDelay(0);
+ scheduler.setDelay(1);
+ scheduler.onInit(consumer);
+ CountDownLatch latch = new CountDownLatch(8);
+ CyclicBarrier barrier = new CyclicBarrier(8);
+ // each concurrent task unschedules the task at the same time
(such as when the repeat count is reached)
+ scheduler.scheduleTask(() -> {
+ try {
+ barrier.await(10, TimeUnit.SECONDS);
+ scheduler.unscheduleTask();
+ } catch (Throwable e) {
+ errors.add(e);
+ } finally {
+ latch.countDown();
+ }
+ });
+ scheduler.start();
+ scheduler.startScheduler();
+ assertTrue(latch.await(10, TimeUnit.SECONDS));
+ scheduler.stop();
+ }
+
+ assertTrue(errors.isEmpty(), errors.toString());
+ }
+
+ @Test
+ public void testPollRunningWhenStoppedDoesNotUpdateState() throws
Exception {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"),
null);
+ consumer.setStartScheduler(false);
+ consumer.start();
+
+ consumer.pollStarted = new CountDownLatch(1);
+ consumer.pollContinue = new CountDownLatch(1);
+ Thread thread = new Thread(consumer::run);
+ thread.start();
+ assertTrue(consumer.pollStarted.await(10, TimeUnit.SECONDS));
+
+ // stop while polling, and then let the poll complete
+ consumer.stop();
+ consumer.pollContinue.countDown();
+ thread.join(10000);
+
+ assertFalse(consumer.isFirstPollDone());
+ assertEquals(0, consumer.getSuccessCounter());
+ }
+
+ @Test
+ public void testErrorThresholdWithoutBackoffMultiplier() {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new
IllegalStateException("Forced"));
+ consumer.setBackoffErrorThreshold(2);
+ consumer.start();
+ try {
+ for (int i = 0; i < 5; i++) {
+ consumer.run();
+ }
+ assertEquals(5, consumer.getErrorCounter());
+ } finally {
+ consumer.stop();
+ }
+ }
+
+ @Test
+ public void testErrorCounterKeptWhenBackoffFinished() {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new
IllegalStateException("Forced"));
+ consumer.setBackoffMultiplier(2);
+ consumer.setBackoffErrorThreshold(1);
+ consumer.start();
+ try {
+ // error
+ consumer.run();
+ assertEquals(1, consumer.getErrorCounter());
+ // backoff (skip)
+ consumer.run();
+ assertEquals(1, consumer.getErrorCounter());
+ // backoff finished, poll again which fails again
+ consumer.run();
+ assertEquals(2, consumer.getErrorCounter());
+ // and backoff again (skip)
+ consumer.run();
+ assertEquals(2, consumer.getErrorCounter());
+ } finally {
+ consumer.stop();
+ }
+ }
+
+ @Test
+ public void testLastErrorDetails() {
+ MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new
MyHttpException());
+ consumer.start();
+ consumer.run();
+ assertEquals(503,
consumer.lastErrorDetails().get(HealthCheck.HTTP_RESPONSE_CODE));
+
+ // another error without http details
+ consumer.setExceptionToThrowOnPoll(new
IllegalStateException("Forced"));
+ consumer.run();
+ assertTrue(consumer.lastError() instanceof IllegalStateException);
+ assertNull(consumer.lastErrorDetails());
+
+ // the error is cleared when stopped
+ consumer.stop();
+ assertNull(consumer.lastError());
+ assertNull(consumer.lastErrorDetails());
+ }
+}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
index bf6cefeeb986..75f97ad96cdf 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
@@ -16,9 +16,9 @@
*/
package org.apache.camel.support;
-import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
+import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
@@ -45,7 +45,8 @@ public class DefaultScheduledPollConsumerScheduler extends
ServiceSupport implem
private Consumer consumer;
private ScheduledExecutorService scheduledExecutorService;
private boolean shutdownExecutor;
- private final List<ScheduledFuture<?>> futures = new ArrayList<>();
+ // thread-safe as the task may be unscheduled by the (concurrent) polling
threads (such as with repeatCount)
+ private final List<ScheduledFuture<?>> futures = new
CopyOnWriteArrayList<>();
private Runnable task;
private int poolSize = 1;
private int concurrentConsumers = 1;
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
b/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
index 8c765e9313b0..f37151f5972b 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
@@ -81,6 +81,12 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
private volatile Map<String, Object> lastErrorDetails;
private final AtomicLong counter = new AtomicLong();
private volatile boolean firstPollDone;
+ // the error counter when the last backoff finished (the error threshold
of the backoff counts from this)
+ private final AtomicLong backoffErrorBase = new AtomicLong();
+ // incremented when the consumer is stopped, so a poll that is still
running does not update the state afterwards
+ private final AtomicLong generation = new AtomicLong();
+ // whether this consumer created the scheduler (and therefore configures
it)
+ private boolean consumerCreatedScheduler;
private volatile boolean forceReady;
protected ScheduledPollConsumer(Endpoint endpoint, Processor processor) {
@@ -143,6 +149,7 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
}
private void doRun() {
+ final long currentGeneration = generation.get();
if (isSuspended()) {
LOG.trace("Cannot start to poll: {} as its suspended",
this.getEndpoint());
return;
@@ -151,8 +158,9 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
// should we backoff if its enabled, and either the idle or error
counter is > the threshold
if (backoffMultiplier > 0
// either idle or error threshold could be not in use, so
check for that and use MAX_VALUE if not in use
- && idleCounter.longValue() >= (backoffIdleThreshold > 0 ?
backoffIdleThreshold : Integer.MAX_VALUE)
- || errorCounter.longValue() >= (backoffErrorThreshold > 0 ?
backoffErrorThreshold : Integer.MAX_VALUE)) {
+ && (idleCounter.longValue() >= (backoffIdleThreshold > 0 ?
backoffIdleThreshold : Integer.MAX_VALUE)
+ || errorCounter.longValue() -
backoffErrorBase.longValue() >= (backoffErrorThreshold > 0
+ ? backoffErrorThreshold : Integer.MAX_VALUE)))
{
final int currentBackoffCounter = backoffCounter.incrementAndGet();
if (currentBackoffCounter < backoffMultiplier) {
// yes we should backoff
@@ -165,9 +173,10 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
}
return;
} else {
- // we are finished with backoff so reset counters
+ // we are finished with backoff so reset counters, but keep
the error counter (the consumer is still
+ // failing until the next poll succeeds) and count the error
threshold from now on
idleCounter.set(0);
- errorCounter.set(0);
+ backoffErrorBase.set(errorCounter.longValue());
backoffCounter.set(0);
successCounter.set(0);
LOG.trace("doRun() backoff finished, resetting counters.");
@@ -220,14 +229,17 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
retryCounter = -1;
LOG.trace("Greedy polling after processing {}
messages", polledMessages);
- // clear any error that might be since we have
successfully polled, otherwise readiness checks might believe the
- // consumer to be unhealthy
- errorCounter.set(0);
- lastError = null;
- lastErrorDetails = null;
-
- // setting firstPollDone to true if greedy
polling is enabled
- firstPollDone = true;
+ if (currentGeneration == generation.get()) {
+ // clear any error that might be since we
have successfully polled, otherwise readiness checks might believe the
+ // consumer to be unhealthy
+ errorCounter.set(0);
+ backoffErrorBase.set(0);
+ lastError = null;
+ lastErrorDetails = null;
+
+ // setting firstPollDone to true if greedy
polling is enabled
+ firstPollDone = true;
+ }
}
} else {
LOG.debug("Cannot begin polling as pollStrategy
returned false: {}", pollStrategy);
@@ -267,11 +279,19 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
}
}
+ if (currentGeneration != generation.get()) {
+ // the consumer has been stopped while polling, so the state
(which has been reset) must not be updated
+ LOG.trace("doRun() done after the consumer was stopped");
+ return;
+ }
+
if (cause != null) {
idleCounter.set(0);
successCounter.set(0);
errorCounter.incrementAndGet();
lastError = cause;
+ // the details are for this error only
+ lastErrorDetails = null;
// enrich last error with http response code if possible
if (cause instanceof HttpResponseAware httpResponseAware) {
int code = httpResponseAware.getHttpResponseCode();
@@ -287,6 +307,7 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
}
successCounter.incrementAndGet();
errorCounter.set(0);
+ backoffErrorBase.set(0);
lastError = null;
lastErrorDetails = null;
}
@@ -641,15 +662,17 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
boolean newScheduler = false;
if (scheduler == null) {
- DefaultScheduledPollConsumerScheduler scheduler
- = new
DefaultScheduledPollConsumerScheduler(scheduledExecutorService);
- scheduler.setDelay(delay);
- scheduler.setInitialDelay(initialDelay);
- scheduler.setTimeUnit(timeUnit);
- scheduler.setUseFixedDelay(useFixedDelay);
- this.scheduler = scheduler;
+ this.scheduler = new
DefaultScheduledPollConsumerScheduler(scheduledExecutorService);
+ consumerCreatedScheduler = true;
newScheduler = true;
}
+ if (consumerCreatedScheduler && scheduler instanceof
DefaultScheduledPollConsumerScheduler defaultScheduler) {
+ // (re)apply the options from this consumer, which may have been
changed while stopped (such as via JMX)
+ defaultScheduler.setDelay(delay);
+ defaultScheduler.setInitialDelay(initialDelay);
+ defaultScheduler.setTimeUnit(timeUnit);
+ defaultScheduler.setUseFixedDelay(useFixedDelay);
+ }
ObjectHelper.notNull(scheduler, "scheduler", this);
scheduler.setCamelContext(getEndpoint().getCamelContext());
@@ -713,15 +736,21 @@ public abstract class ScheduledPollConsumer extends
DefaultConsumer
ServiceHelper.stopAndShutdownServices(scheduler);
}
+ // a poll that is still running must not update the state after it has
been cleared
+ generation.incrementAndGet();
+
// clear counters
backoffCounter.set(0);
idleCounter.set(0);
errorCounter.set(0);
+ backoffErrorBase.set(0);
successCounter.set(0);
counter.set(0);
// clear ready state
firstPollDone = false;
forceReady = false;
+ lastError = null;
+ lastErrorDetails = null;
super.doStop();
}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 3e50f5aa2cba..8fbb0a4c0b5e 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -3446,6 +3446,15 @@ redelivery now follows
`negativeAckRedeliveryDelayMicros`, which defaults to 60
A route that wants the previous timing can set
`negativeAckRedeliveryDelayMicros=10000000`.
+=== camel-core - scheduled poll consumer
+
+The `delay`, `initialDelay`, `timeUnit` and `useFixedDelay` of a scheduled
poll consumer that are changed while the
+consumer is stopped (such as using JMX) are now used when the consumer is
started again. Previously the values from
+the first start were kept.
+
+When a backoff (`backoffMultiplier`) finishes, the error counter of the
consumer (used by the health checks and JMX)
+is no longer reset to 0. It is reset when a poll succeeds. The number of
skipped polls is the same as before.
+
=== camel-seda - multipleConsumers broadcasts to consumers with different uri
options
With `multipleConsumers=true` every consumer of a SEDA queue now receives a
copy of each message, also when the consumers