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 7ae3b275759d CAMEL-24991: camel-core - Aggregate EIP: fix bugs found
in a deep review (#26847)
7ae3b275759d is described below
commit 7ae3b275759d437d4f98dd1e4414548ca0d28f18
Author: Claus Ibsen <[email protected]>
AuthorDate: Sun Sep 27 10:10:02 2026 +0200
CAMEL-24991: camel-core - Aggregate EIP: fix bugs found in a deep review
(#26847)
- With optimistic locking, force completion and force discarding of
groups stopped at the first group another Camel instance had
completed, and the exception reached the caller. It is now handled
per group, as the completion interval task already does.
- Force discarding a group only discarded it with
discardOnAggregationFailure enabled. Otherwise the group was removed
without being confirmed, so a recoverable repository could send it
later.
- The optimistic locking retry delay is capped at the documented
default of 1000 ms, and the exponential back-off no longer overflows.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../processor/aggregate/AggregateProcessor.java | 61 +++++++++++---
.../aggregate/OptimisticLockRetryPolicy.java | 12 ++-
.../AggregateForceDiscardingConfirmTest.java | 74 +++++++++++++++++
...regateOptimisticLockingForceCompletionTest.java | 96 ++++++++++++++++++++++
.../aggregator/OptimisticLockRetryPolicyTest.java | 24 ++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 8 ++
6 files changed, 260 insertions(+), 15 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
index 9bb85dc6f065..aa569e1a4a5a 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
@@ -1839,13 +1839,11 @@ public class AggregateProcessor extends
BaseProcessorSupport
try {
Exchange exchange = aggregationRepository.get(camelContext, key);
if (exchange != null) {
- total = 1;
LOG.trace("Force completion triggered for correlation key:
{}", key);
// indicate it was completed by a force completion request
exchange.setProperty(ExchangePropertyKey.AGGREGATED_COMPLETED_BY,
COMPLETED_BY_FORCE);
- Exchange answer = onCompletion(key, exchange, exchange, false,
false);
- if (answer != null) {
- onSubmitCompletion(key, answer);
+ if (forceCompletion(key, exchange)) {
+ total = 1;
}
}
} finally {
@@ -1886,10 +1884,7 @@ public class AggregateProcessor extends
BaseProcessorSupport
LOG.trace("Force completion triggered for correlation
key: {}", key);
// indicate it was completed by a force completion
request
exchange.setProperty(ExchangePropertyKey.AGGREGATED_COMPLETED_BY,
COMPLETED_BY_FORCE);
- Exchange answer = onCompletion(key, exchange,
exchange, false, false);
- if (answer != null) {
- onSubmitCompletion(key, answer);
- }
+ forceCompletion(key, exchange);
}
}
} finally {
@@ -1912,10 +1907,10 @@ public class AggregateProcessor extends
BaseProcessorSupport
try {
Exchange exchange = aggregationRepository.get(camelContext, key);
if (exchange != null) {
- total = 1;
LOG.trace("Force discarded triggered for correlation key: {}",
key);
- // force discarding by setting aggregate failed as true
- onCompletion(key, exchange, exchange, false, true);
+ if (forceDiscarding(key, exchange)) {
+ total = 1;
+ }
}
} finally {
lock.unlock();
@@ -1953,8 +1948,7 @@ public class AggregateProcessor extends
BaseProcessorSupport
Exchange exchange =
aggregationRepository.get(camelContext, key);
if (exchange != null) {
LOG.trace("Force discarded triggered for correlation
key: {}", key);
- // force discarding by setting aggregate failed as true
- onCompletion(key, exchange, exchange, false, true);
+ forceDiscarding(key, exchange);
}
}
} finally {
@@ -1969,6 +1963,47 @@ public class AggregateProcessor extends
BaseProcessorSupport
return total;
}
+ /**
+ * Completes the group and sends the aggregated exchange. Must be called
while holding the lock.
+ *
+ * @return {@code false} if another Camel instance has already completed
the group (optimistic locking)
+ */
+ private boolean forceCompletion(String key, Exchange exchange) {
+ try {
+ Exchange answer = onCompletion(key, exchange, exchange, false,
false);
+ if (answer != null) {
+ onSubmitCompletion(key, answer);
+ }
+ return true;
+ } catch
(OptimisticLockingAggregationRepository.OptimisticLockingException e) {
+ LOG.debug("Another Camel instance has already completed the group
with correlation key: {}", key);
+ return false;
+ }
+ }
+
+ /**
+ * Discards the group. Must be called while holding the lock.
+ *
+ * @return {@code false} if another Camel instance has already removed the
group from the repository (optimistic
+ * locking)
+ */
+ private boolean forceDiscarding(String key, Exchange exchange) {
+ try {
+ // force discarding by setting aggregate failed as true
+ Exchange answer = onCompletion(key, exchange, exchange, false,
true);
+ if (answer != null) {
+ // onCompletion only discards on aggregation failure when
discardOnAggregationFailure is enabled,
+ // so discard here, as otherwise the group is removed without
being confirmed (and a recoverable
+ // repository would recover and send it later)
+ discard(key, answer);
+ }
+ return true;
+ } catch
(OptimisticLockingAggregationRepository.OptimisticLockingException e) {
+ LOG.debug("Another Camel instance has already completed the group
with correlation key: {}", key);
+ return false;
+ }
+ }
+
/**
* Synchronization class to avoid busy-loop when waiting for exchanges to
be processed during shutdown.
*/
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/OptimisticLockRetryPolicy.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/OptimisticLockRetryPolicy.java
index 4f1f8b95c381..24a0032c9196 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/OptimisticLockRetryPolicy.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/OptimisticLockRetryPolicy.java
@@ -46,7 +46,7 @@ public class OptimisticLockRetryPolicy {
private int maximumRetries;
private long retryDelay = 50L;
- private long maximumRetryDelay;
+ private long maximumRetryDelay = DEFAULT_MAXIMUM_RETRY_DELAY;
private boolean exponentialBackOff = true;
private boolean randomBackOff;
@@ -68,7 +68,7 @@ public class OptimisticLockRetryPolicy {
long sleepFor = 0;
if (retryDelay > 0 || randomBackOff) {
sleepFor = exponentialBackOff
- ? (retryDelay << retryCounter)
+ ? exponentialDelay(retryCounter)
: (randomBackOff
? ThreadLocalRandom.current() // NOSONAR
.nextInt((int) (maximumRetryDelay > 0 ?
maximumRetryDelay : DEFAULT_MAXIMUM_RETRY_DELAY))
@@ -80,6 +80,14 @@ public class OptimisticLockRetryPolicy {
return sleepFor;
}
+ private long exponentialDelay(int retryCounter) {
+ // the retry counter keeps growing when retrying forever, so do not
let the shift overflow
+ if (retryCounter >= Long.numberOfLeadingZeros(retryDelay) - 1) {
+ return Long.MAX_VALUE;
+ }
+ return retryDelay << retryCounter;
+ }
+
public int getMaximumRetries() {
return maximumRetries;
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateForceDiscardingConfirmTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateForceDiscardingConfirmTest.java
new file mode 100644
index 000000000000..865f9f62d0da
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateForceDiscardingConfirmTest.java
@@ -0,0 +1,74 @@
+/*
+ * 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.processor.aggregator;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.processor.aggregate.AggregateController;
+import org.apache.camel.processor.aggregate.DefaultAggregateController;
+import org.apache.camel.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.processor.aggregate.StringAggregationStrategy;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * Force discarding a group discards it (confirms it in the repository), also
when discardOnAggregationFailure is not
+ * enabled.
+ */
+public class AggregateForceDiscardingConfirmTest extends ContextTestSupport {
+
+ private final AggregateController controller = new
DefaultAggregateController();
+ private final List<String> confirmed = new ArrayList<>();
+
+ @Test
+ public void testForceDiscardingOfGroupConfirms() throws Exception {
+ getMockEndpoint("mock:aggregated").expectedMessageCount(0);
+
+ template.sendBodyAndHeader("direct:start", "A", "id", "A");
+ template.sendBodyAndHeader("direct:start", "B", "id", "B");
+
+ assertEquals(1, controller.forceDiscardingOfGroup("A"));
+ assertEquals(1, controller.forceDiscardingOfAllGroups());
+
+ assertMockEndpointsSatisfied();
+ assertEquals(2, confirmed.size(), "the discarded groups should be
confirmed");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), new
StringAggregationStrategy()).completionSize(10)
+ .aggregationRepository(new
MemoryAggregationRepository() {
+ @Override
+ public void confirm(CamelContext camelContext,
String exchangeId) {
+ confirmed.add(exchangeId);
+ }
+ })
+ .aggregateController(controller)
+ .to("mock:aggregated");
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingForceCompletionTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingForceCompletionTest.java
new file mode 100644
index 000000000000..b91c0adeff0f
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingForceCompletionTest.java
@@ -0,0 +1,96 @@
+/*
+ * 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.processor.aggregator;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.processor.aggregate.AggregateController;
+import org.apache.camel.processor.aggregate.DefaultAggregateController;
+import org.apache.camel.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.processor.aggregate.StringAggregationStrategy;
+import org.apache.camel.spi.OptimisticLockingAggregationRepository;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * With optimistic locking, a group that another node completed first does not
stop the force completion of the other
+ * groups.
+ */
+public class AggregateOptimisticLockingForceCompletionTest extends
ContextTestSupport {
+
+ private final AggregateController controller = new
DefaultAggregateController();
+
+ @Test
+ public void testForceCompletionContinuesAfterConflict() throws Exception {
+
getMockEndpoint("mock:aggregated").expectedBodiesReceivedInAnyOrder("A", "C");
+
+ template.sendBodyAndHeader("direct:start", "A", "id", "A");
+ template.sendBodyAndHeader("direct:start", "B", "id", "B");
+ template.sendBodyAndHeader("direct:start", "C", "id", "C");
+
+ controller.forceCompletionOfAllGroups();
+
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testForceDiscardingContinuesAfterConflict() throws Exception {
+ getMockEndpoint("mock:aggregated").expectedMessageCount(0);
+
+ template.sendBodyAndHeader("direct:start", "A", "id", "A");
+ template.sendBodyAndHeader("direct:start", "B", "id", "B");
+ template.sendBodyAndHeader("direct:start", "C", "id", "C");
+
+ // only group B is left (completed by another node, which this test
repository never removes)
+ controller.forceDiscardingOfAllGroups();
+ assertEquals(1, controller.forceDiscardingOfAllGroups());
+
+ assertMockEndpointsSatisfied();
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), new
StringAggregationStrategy()).completionSize(10)
+ .aggregationRepository(new
StolenRepository()).optimisticLocking()
+ .aggregateController(controller)
+ .to("mock:aggregated");
+ }
+ };
+ }
+
+ // another node has completed group B first
+ private static class StolenRepository extends MemoryAggregationRepository {
+ StolenRepository() {
+ super(true);
+ }
+
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ if ("B".equals(key)) {
+ throw new
OptimisticLockingAggregationRepository.OptimisticLockingException();
+ }
+ super.remove(camelContext, key, exchange);
+ }
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/OptimisticLockRetryPolicyTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/OptimisticLockRetryPolicyTest.java
index 66271cc53af2..4dadf222d8ba 100644
---
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/OptimisticLockRetryPolicyTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/OptimisticLockRetryPolicyTest.java
@@ -105,6 +105,30 @@ class OptimisticLockRetryPolicyTest {
}
}
+ @Test
+ void testDefaultMaximumRetryDelay() {
+ // the default maximum retry delay is 1 second, as documented
+ OptimisticLockRetryPolicy policy = new OptimisticLockRetryPolicy();
+ assertEquals(1000L, policy.getMaximumRetryDelay());
+ assertDelay(100L, policy.getDelay(1));
+ assertDelay(1000L, policy.getDelay(10));
+ assertDelay(1000L, policy.getDelay(100));
+ }
+
+ @Test
+ void testExponentialBackOffDoesNotOverflow() {
+ OptimisticLockRetryPolicy policy = new OptimisticLockRetryPolicy();
+ policy.setRetryDelay(50L);
+ policy.setMaximumRetryDelay(0L);
+
+ long previous = 0;
+ for (int i = 0; i < 200; i++) {
+ long delay = policy.getDelay(i);
+ assertTrue(delay >= previous, "delay should not decrease at retry
" + i);
+ previous = delay;
+ }
+ }
+
private long getDelay(OptimisticLockRetryPolicy policy, int i) {
return policy.getDelay(i);
}
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 411e8a2d67f7..5f750d4fb38e 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
@@ -114,6 +114,14 @@ This can break an application that has such an invalid
value configured and work
good enough, as it now fails when starting. To keep the previous result,
correct the configured value, or remove the
property so the `defaultValue` is used.
+=== Aggregate EIP - optimistic locking retry delay
+
+When optimistic locking is enabled without configuring an
`optimisticLockRetryPolicy`, the retry delay between
+attempts is now capped at 1 second, which is the documented default of
`maximumRetryDelay`. Prior to Camel 4.23 the
+delay was not capped in that case and doubled on every attempt, so after a few
failed attempts an exchange could wait
+minutes or hours before trying again. To keep an uncapped delay, configure an
`optimisticLockRetryPolicy` with
+`maximumRetryDelay` set to `0`.
+
=== Circuit Breaker EIP
The exchange property `CamelCircuitBreakerResponseRejected` is now also set
inside the `onFallback`,