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 ebecd1afbb45 CAMEL-25211: camel-hazelcast - recover the completed
aggregation with the message that completed it (#27160)
ebecd1afbb45 is described below
commit ebecd1afbb45233aae4658aa4733b64743cd561b
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 12:24:52 2026 +0530
CAMEL-25211: camel-hazelcast - recover the completed aggregation with the
message that completed it (#27160)
With recovery enabled (the default) and without optimistic locking,
HazelcastAggregationRepository.remove() stored the entry it removed from
the map as the exchange to recover. When a group is completed by an
incoming exchange, the Aggregate EIP does not add the aggregated exchange
to the repository before it removes the group, so that entry is the group
before the last message. A completed exchange whose processing failed was
recovered without the message that completed the group. Store the
marshalled exchange passed to remove(), as the optimistic branch already
does, and as CAMEL-24946 did for KeyValueAggregationRepository.
ReplicatedHazelcastAggregationRepository keeps its groups and completed
exchanges in ReplicatedMaps, but its non-optimistic remove() with recovery
ran a Hazelcast transaction on the IMaps of the same names. The group was
never found there, so the removed value was null, the put of the completed
exchange failed with "value can't be null", and the transaction was rolled
back: every group of two or more messages failed to complete and stayed in
the repository. A ReplicatedMap cannot take part in a transaction, so
remove() now stores the given exchange in the replicated completed map and
then removes the group, under the per-key lock that add() uses.
HazelcastAggregationRepositoryRecoverTest sends a group of three messages
whose processing fails once: the recovered exchange must hold all three
messages, for both repositories.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../hazelcast/HazelcastAggregationRepository.java | 6 +-
.../ReplicatedHazelcastAggregationRepository.java | 45 +++-------
.../HazelcastAggregationRepositoryRecoverTest.java | 98 ++++++++++++++++++++++
3 files changed, 112 insertions(+), 37 deletions(-)
diff --git
a/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepository.java
b/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepository.java
index b08cc9643a95..5a836a754a4b 100644
---
a/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepository.java
+++
b/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepository.java
@@ -390,10 +390,12 @@ public class HazelcastAggregationRepository extends
ServiceSupport
TransactionalMap<String, DefaultExchangeHolder> tCache =
tCtx.getMap(cache.getName());
TransactionalMap<String, DefaultExchangeHolder>
tPersistentCache = tCtx.getMap(persistedCache.getName());
- DefaultExchangeHolder removedHolder = tCache.remove(key);
+ tCache.remove(key);
LOG.trace("Putting an exchange with ID {} for key {} into
a recoverable storage in a thread-safe manner.",
exchange.getExchangeId(), key);
- tPersistentCache.put(exchange.getExchangeId(),
removedHolder);
+ // store the given exchange and not the removed entry:
when a group is completed by an incoming
+ // exchange, the last aggregated exchange is not added to
the repository before it is removed
+ tPersistentCache.put(exchange.getExchangeId(), holder);
tCtx.commitTransaction();
LOG.trace("Removed an exchange with ID {} for key {} in a
thread-safe manner.", exchange.getExchangeId(),
diff --git
a/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/ReplicatedHazelcastAggregationRepository.java
b/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/ReplicatedHazelcastAggregationRepository.java
index a8c4b1e4dd9f..5941f60849d2 100644
---
a/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/ReplicatedHazelcastAggregationRepository.java
+++
b/components/camel-hazelcast/src/main/java/org/apache/camel/processor/aggregate/hazelcast/ReplicatedHazelcastAggregationRepository.java
@@ -25,12 +25,8 @@ import com.hazelcast.config.XmlConfigBuilder;
import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.map.IMap;
-import com.hazelcast.transaction.TransactionContext;
-import com.hazelcast.transaction.TransactionOptions;
-import com.hazelcast.transaction.TransactionalMap;
import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
-import org.apache.camel.RuntimeCamelException;
import org.apache.camel.component.hazelcast.HazelcastSerializationFilterHelper;
import org.apache.camel.spi.OptimisticLockingAggregationRepository;
import org.apache.camel.spi.RecoverableAggregationRepository;
@@ -278,39 +274,18 @@ public class ReplicatedHazelcastAggregationRepository
extends HazelcastAggregati
} else {
if (useRecovery) {
LOG.trace("Removing an exchange with ID {} for key {} in a
thread-safe manner.", exchange.getExchangeId(), key);
- // The only considerable case for transaction usage is fault
tolerance:
- // the transaction will be rolled back automatically (default
timeout is 2 minutes)
- // if no commit occurs during the timeout. So we are still
consistent whether local node crashes.
- TransactionOptions tOpts = new TransactionOptions();
-
-
tOpts.setTransactionType(TransactionOptions.TransactionType.ONE_PHASE);
- TransactionContext tCtx =
hazelcastInstance.newTransactionContext(tOpts);
-
+ // a ReplicatedMap cannot take part in a Hazelcast
transaction, so the completed exchange is stored for
+ // recovery before the group is removed, under the lock that
add() uses for the key.
+ // The given exchange is stored and not the removed entry:
when a group is completed by an incoming
+ // exchange, the last aggregated exchange is not added to the
repository before it is removed
+ lockMap.lock(key);
try {
- tCtx.beginTransaction();
-
- TransactionalMap<String, DefaultExchangeHolder> tCache =
tCtx.getMap(mapName);
- TransactionalMap<String, DefaultExchangeHolder>
tPersistentCache = tCtx.getMap(persistenceMapName);
-
- DefaultExchangeHolder removedHolder = tCache.remove(key);
- LOG.trace("Putting an exchange with ID {} for key {} into
a recoverable storage in a thread-safe manner.",
- exchange.getExchangeId(), key);
- tPersistentCache.put(exchange.getExchangeId(),
removedHolder);
-
- tCtx.commitTransaction();
- LOG.trace("Removed an exchange with ID {} for key {} in a
thread-safe manner.", exchange.getExchangeId(),
- key);
- LOG.trace("Put an exchange with ID {} for key {} into a
recoverable storage in a thread-safe manner.",
- exchange.getExchangeId(), key);
- } catch (Exception throwable) {
- tCtx.rollbackTransaction();
-
- final String msg = String.format(
- "Transaction with ID %s was rolled back for remove
operation with a key %s and an Exchange ID %s.",
- tCtx.getTxnId(), key, exchange.getExchangeId());
- LOG.warn(msg, throwable);
- throw new RuntimeCamelException(msg, throwable);
+ replicatedPersistedCache.put(exchange.getExchangeId(),
holder);
+ replicatedCache.remove(key);
+ } finally {
+ lockMap.unlock(key);
}
+ LOG.trace("Removed an exchange with ID {} for key {} in a
thread-safe manner.", exchange.getExchangeId(), key);
} else {
replicatedCache.remove(key);
}
diff --git
a/components/camel-hazelcast/src/test/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepositoryRecoverTest.java
b/components/camel-hazelcast/src/test/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepositoryRecoverTest.java
new file mode 100644
index 000000000000..65199be9b167
--- /dev/null
+++
b/components/camel-hazelcast/src/test/java/org/apache/camel/processor/aggregate/hazelcast/HazelcastAggregationRepositoryRecoverTest.java
@@ -0,0 +1,98 @@
+/*
+ * 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.aggregate.hazelcast;
+
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A completed aggregation whose processing fails is recovered with all its
messages, including the one that completed
+ * the group (which is never added to the repository).
+ */
+public class HazelcastAggregationRepositoryRecoverTest extends
HazelcastAggregationRepositoryCamelTestSupport {
+
+ @Test
+ public void testRecoveredExchangeContainsTheLastMessage() throws Exception
{
+ HazelcastAggregationRepository repository
+ = new HazelcastAggregationRepository("recoverRepo", false,
getFirstInstance());
+
+ assertRecoveredWithAllMessages(repository, "imap");
+ }
+
+ @Test
+ public void testReplicatedRecoveredExchangeContainsTheLastMessage() throws
Exception {
+ ReplicatedHazelcastAggregationRepository repository
+ = new
ReplicatedHazelcastAggregationRepository("replicatedRecoverRepo", false,
getFirstInstance());
+
+ assertRecoveredWithAllMessages(repository, "replicated");
+ }
+
+ private void assertRecoveredWithAllMessages(HazelcastAggregationRepository
repository, String name) throws Exception {
+ repository.setRecoveryInterval(100);
+ AtomicInteger failures = new AtomicInteger();
+
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:" + name)
+ .aggregate(header("id"), (oldExchange, newExchange) ->
{
+ if (oldExchange == null) {
+ return newExchange;
+ }
+ String body =
oldExchange.getIn().getBody(String.class) + "+"
+ +
newExchange.getIn().getBody(String.class);
+ oldExchange.getIn().setBody(body);
+ return oldExchange;
+ })
+ .aggregationRepository(repository)
+ .completionSize(3)
+ .to("mock:" + name)
+ .process(exchange -> {
+ if (failures.getAndIncrement() == 0) {
+ throw new IllegalStateException("Forced
failure after the aggregation");
+ }
+ });
+ }
+ });
+
+ MockEndpoint mock = getMockEndpoint("mock:" + name);
+ // the completed aggregation fails once and is then recovered with all
three messages
+ mock.expectedBodiesReceived("a+b+c", "a+b+c");
+
+ template.sendBodyAndHeader("direct:" + name, "a", "id", "group");
+ template.sendBodyAndHeader("direct:" + name, "b", "id", "group");
+ template.sendBodyAndHeader("direct:" + name, "c", "id", "group");
+
+ mock.assertIsSatisfied();
+
assertNull(mock.getReceivedExchanges().get(0).getIn().getHeader(Exchange.REDELIVERED));
+ assertEquals(Boolean.TRUE,
mock.getReceivedExchanges().get(1).getIn().getHeader(Exchange.REDELIVERED));
+ // the recovered exchange is confirmed after its successful
redelivery, so it is not sent again
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
repository.scan(context).isEmpty());
+ assertEquals(2, mock.getReceivedCounter());
+ assertTrue(repository.getKeys().isEmpty());
+ }
+}