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 1de5a7c0cb52 CAMEL-25210: camel-cassandraql - keep completed exchanges
for recovery in the aggregation repository (#27159)
1de5a7c0cb52 is described below
commit 1de5a7c0cb52d87515728f263a7faca2e798b907
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 21:10:55 2026 +0530
CAMEL-25210: camel-cassandraql - keep completed exchanges for recovery in
the aggregation repository (#27159)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../cassandra/CassandraAggregationRepository.java | 88 ++++-----
.../CassandraAggregationRepositoryIT.java | 27 ++-
...CassandraAggregationRepositoryRecoveryTest.java | 202 +++++++++++++++++++++
.../NamedCassandraAggregationRepositoryIT.java | 27 ++-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 18 ++
5 files changed, 313 insertions(+), 49 deletions(-)
diff --git
a/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
b/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
index f3c2b9a7d913..92662ae9d252 100644
---
a/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
+++
b/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
@@ -18,6 +18,7 @@ package org.apache.camel.processor.aggregate.cassandra;
import java.io.IOException;
import java.nio.ByteBuffer;
+import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
@@ -42,7 +43,6 @@ import
org.apache.camel.utils.cassandra.CassandraSessionHolder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker;
import static org.apache.camel.utils.cassandra.CassandraUtils.append;
import static
org.apache.camel.utils.cassandra.CassandraUtils.applyConsistencyLevel;
import static org.apache.camel.utils.cassandra.CassandraUtils.concat;
@@ -64,6 +64,11 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
private static final Logger LOGGER =
LoggerFactory.getLogger(CassandraAggregationRepository.class);
+ /**
+ * Prefix of the aggregation key under which a completed exchange is kept,
by its exchange id, until it is confirmed
+ */
+ private static final String RECOVERY_KEY_PREFIX = "camel-recovery:";
+
private final CassandraCamelCodec exchangeCodec = new
CassandraCamelCodec();
@Metadata(description = "Cassandra session", required = true)
@@ -97,10 +102,6 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
* Prepared statement used to get exchangeIds and exchange ids
*/
private PreparedStatement selectKeyIdStatement;
- /**
- * Prepared statement used to delete with key and exchange id
- */
- private PreparedStatement deleteIfIdStatement;
@Metadata(description = "Sets the interval between recovery scans",
defaultValue = "5000")
private long recoveryInterval = 5000;
@@ -167,7 +168,6 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
initSelectStatement();
initDeleteStatement();
initSelectKeyIdStatement();
- initDeleteIfIdStatement();
}
@Override
@@ -190,13 +190,17 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
*/
@Override
public Exchange add(CamelContext camelContext, String key, Exchange
exchange) {
+ insert(key, exchange);
+ return exchange;
+ }
+
+ private void insert(String key, Exchange exchange) {
final Object[] idValues = getPKValues(key);
LOGGER.debug("Inserting key {} exchange {}", idValues, exchange);
try {
ByteBuffer marshalledExchange =
exchangeCodec.marshallExchange(exchange, allowSerializedHeaders);
Object[] cqlParams = concat(idValues, new Object[] {
exchange.getExchangeId(), marshalledExchange });
getSession().execute(insertStatement.bind(cqlParams));
- return exchange;
} catch (IOException iOException) {
throw new CassandraAggregationException("Failed to write
exchange", exchange, iOException);
}
@@ -236,29 +240,16 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
//
-------------------------------------------------------------------------
// Confirm exchange in repository
- private void initDeleteIfIdStatement() {
- Delete delete = generateDelete(table, pkColumns, false);
- Delete deleteIf =
delete.ifColumn(exchangeIdColumn).isEqualTo(bindMarker());
- SimpleStatement statement = applyConsistencyLevel(deleteIf.build(),
writeConsistencyLevel);
- LOGGER.debug("Generated Delete If Id {}", statement);
- deleteIfIdStatement = getSession().prepare(statement);
- }
/**
- * Remove exchange by Id from aggregation table.
+ * Remove the completed exchange kept for recovery.
*/
@Override
public void confirm(CamelContext camelContext, String exchangeId) {
- String keyColumn = getKeyColumn();
- LOGGER.debug("Selecting Ids");
- List<Row> rows = selectKeyIds();
- for (Row row : rows) {
- if (row.getString(exchangeIdColumn).equals(exchangeId)) {
- String key = row.getString(keyColumn);
- Object[] cqlParams = append(getPKValues(key), exchangeId);
- LOGGER.debug("Deleting If Id {}", cqlParams);
- getSession().execute(deleteIfIdStatement.bind(cqlParams));
- }
+ if (useRecovery) {
+ Object[] idValues = getPKValues(recoveryKey(exchangeId));
+ LOGGER.debug("Deleting key {}", (Object) idValues);
+ getSession().execute(deleteStatement.bind(idValues));
}
}
@@ -277,6 +268,13 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
*/
@Override
public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ if (useRecovery) {
+ // the aggregation is complete but the exchange has not been
processed yet, so keep it under its exchange
+ // id until it is confirmed, so the recover task can pick it up.
This is the given exchange, as the stored
+ // one does not contain the exchange that completed the
aggregation. It is written before the aggregation
+ // is deleted, so a failure in between does not lose it.
+ insert(recoveryKey(exchange.getExchangeId()), exchange);
+ }
Object[] idValues = getPKValues(key);
LOGGER.debug("Deleting key {}", (Object) idValues);
getSession().execute(deleteStatement.bind(idValues));
@@ -304,7 +302,7 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
}
/**
- * Get aggregation exchangeIds from aggregation table.
+ * Get the aggregation keys of the aggregations in progress from
aggregation table.
*/
@Override
public Set<String> getKeys() {
@@ -312,7 +310,10 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
Set<String> keys = new HashSet<>(rows.size());
String keyColumnName = getKeyColumn();
for (Row row : rows) {
- keys.add(row.getString(keyColumnName));
+ String key = row.getString(keyColumnName);
+ if (!isRecoveryKey(key)) {
+ keys.add(key);
+ }
}
return keys;
}
@@ -324,30 +325,35 @@ public class CassandraAggregationRepository extends
ServiceSupport implements Re
*/
@Override
public Set<String> scan(CamelContext camelContext) {
+ if (!useRecovery) {
+ return Collections.emptySet();
+ }
List<Row> rows = selectKeyIds();
- Set<String> exchangeIds = new HashSet<>(rows.size());
+ Set<String> exchangeIds = new HashSet<>();
+ String keyColumnName = getKeyColumn();
for (Row row : rows) {
- exchangeIds.add(row.getString(exchangeIdColumn));
+ String key = row.getString(keyColumnName);
+ if (isRecoveryKey(key)) {
+ exchangeIds.add(key.substring(RECOVERY_KEY_PREFIX.length()));
+ }
}
return exchangeIds;
}
/**
- * Get exchange by exchange ID. This is far from optimal.
+ * Get the completed exchange kept for recovery by exchange ID.
*/
@Override
public Exchange recover(CamelContext camelContext, String exchangeId) {
- List<Row> rows = selectKeyIds();
- String keyColumnName = getKeyColumn();
- String lKey = null;
- for (Row row : rows) {
- String lExchangeId = row.getString(exchangeIdColumn);
- if (lExchangeId.equals(exchangeId)) {
- lKey = row.getString(keyColumnName);
- break;
- }
- }
- return lKey == null ? null : get(camelContext, lKey);
+ return useRecovery ? get(camelContext, recoveryKey(exchangeId)) : null;
+ }
+
+ private static String recoveryKey(String exchangeId) {
+ return RECOVERY_KEY_PREFIX + exchangeId;
+ }
+
+ private static boolean isRecoveryKey(String key) {
+ return key != null && key.startsWith(RECOVERY_KEY_PREFIX);
}
//
-------------------------------------------------------------------------
diff --git
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
index f08e56bffed0..d2a47bb35dcb 100644
---
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
+++
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
@@ -144,11 +144,16 @@ public class CassandraAggregationRepositoryIT extends
BaseCassandra {
aggregationRepository.add(context, key, exchange);
assertTrue(exists(key));
}
+ Exchange completed = new DefaultExchange(context);
+ completed.setExchangeId("Exchange_2");
+ aggregationRepository.remove(context, "Confirm_2", completed);
+ assertFalse(exists("Confirm_2"));
+ assertTrue(exists("camel-recovery:Exchange_2"));
// When
aggregationRepository.confirm(context, "Exchange_2");
// Then
+ assertFalse(exists("camel-recovery:Exchange_2"));
assertTrue(exists("Confirm_1"));
- assertFalse(exists("Confirm_2"));
assertTrue(exists("Confirm_3"));
}
@@ -171,6 +176,12 @@ public class CassandraAggregationRepositoryIT extends
BaseCassandra {
}
}
+ private void removeExchange(String key) {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.setExchangeId("Exchange-" + key);
+ aggregationRepository.remove(context, key, exchange);
+ }
+
private void addExchanges(String... keys) {
for (String key : keys) {
Exchange exchange = new DefaultExchange(context);
@@ -184,12 +195,16 @@ public class CassandraAggregationRepositoryIT extends
BaseCassandra {
// Given
String[] keys = { "Scan1", "Scan2" };
addExchanges(keys);
+ // aggregations in progress are not to be recovered
+
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan1"));
+
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan2"));
// When
+ removeExchange("Scan2");
Set<String> exchangeIdSet = aggregationRepository.scan(context);
// Then
- for (String key : keys) {
- assertTrue(exchangeIdSet.contains("Exchange-" + key));
- }
+ assertTrue(exchangeIdSet.contains("Exchange-Scan2"));
+ assertFalse(exchangeIdSet.contains("Exchange-Scan1"));
+
assertFalse(aggregationRepository.getKeys().contains("camel-recovery:Exchange-Scan2"));
}
@Test
@@ -197,11 +212,15 @@ public class CassandraAggregationRepositoryIT extends
BaseCassandra {
// Given
String[] keys = { "Recover1", "Recover2" };
addExchanges(keys);
+ removeExchange("Recover2");
// When
+ Exchange exchange1 = aggregationRepository.recover(context,
"Exchange-Recover1");
Exchange exchange2 = aggregationRepository.recover(context,
"Exchange-Recover2");
Exchange exchange3 = aggregationRepository.recover(context,
"Exchange-Recover3");
// Then
+ assertNull(exchange1);
assertNotNull(exchange2);
+ assertEquals("Exchange-Recover2", exchange2.getExchangeId());
assertNull(exchange3);
}
diff --git
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
new file mode 100644
index 000000000000..edd5a6fdea7f
--- /dev/null
+++
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
@@ -0,0 +1,202 @@
+/*
+ * 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.cassandra;
+
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.cql.BoundStatement;
+import com.datastax.oss.driver.api.core.cql.PreparedStatement;
+import com.datastax.oss.driver.api.core.cql.ResultSet;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+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;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * The recover task of the Aggregate EIP must only see completed exchanges
that were not confirmed: an aggregation that
+ * is still open is not recovered, and a completed aggregation whose
processing fails is recovered with all its
+ * messages. The repository runs against an in-memory table behind a mocked
{@link CqlSession}.
+ */
+public class CassandraAggregationRepositoryRecoveryTest extends
CamelTestSupport {
+
+ private final AtomicInteger scans = new AtomicInteger();
+
+ @Test
+ public void testOpenAggregationIsNotRecovered() throws Exception {
+ CassandraAggregationRepository repository = createRepository();
+ addRoute(repository, "open", 3);
+
+ MockEndpoint mock = getMockEndpoint("mock:open");
+ mock.expectedMessageCount(0);
+
+ template.sendBodyAndHeader("direct:open", "a", "id", "group");
+ template.sendBodyAndHeader("direct:open", "b", "id", "group");
+
+ // let the recover task run a few times while the aggregation is open
+ int scanned = scans.get();
+ await().atMost(10, TimeUnit.SECONDS).until(() -> scans.get() >=
scanned + 3);
+
+ mock.assertIsSatisfied();
+ assertEquals("a+b", repository.get(context,
"group").getIn().getBody(String.class));
+ }
+
+ @Test
+ public void testCompletedAggregationIsRecoveredWithAllMessages() throws
Exception {
+ CassandraAggregationRepository repository = createRepository();
+ addRoute(repository, "completed", 3);
+
+ MockEndpoint mock = getMockEndpoint("mock:completed");
+ // the completed aggregation fails once and is then recovered with all
three messages
+ mock.expectedBodiesReceived("a+b+c", "a+b+c");
+
+ template.sendBodyAndHeader("direct:completed", "a", "id", "group");
+ template.sendBodyAndHeader("direct:completed", "b", "id", "group");
+ template.sendBodyAndHeader("direct:completed", "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));
+ // confirmed after the successful redelivery
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
repository.scan(context).isEmpty());
+ assertTrue(repository.getKeys().isEmpty());
+ }
+
+ private CassandraAggregationRepository createRepository() {
+ CassandraAggregationRepository repository
+ = new CassandraAggregationRepository(inMemorySession(new
ConcurrentHashMap<>())) {
+ @Override
+ public Set<String> scan(CamelContext camelContext) {
+ try {
+ return super.scan(camelContext);
+ } finally {
+ scans.incrementAndGet();
+ }
+ }
+ };
+ repository.setRecoveryInterval(100);
+ return repository;
+ }
+
+ private void addRoute(CassandraAggregationRepository repository, String
name, int completionSize) throws Exception {
+ 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(completionSize)
+ .to("mock:" + name)
+ .process(exchange -> {
+ if (failures.getAndIncrement() == 0) {
+ throw new IllegalStateException("Forced
failure after the aggregation");
+ }
+ });
+ }
+ });
+ }
+
+ /**
+ * A session backed by a map of KEY to (KEY, EXCHANGE_ID, EXCHANGE), which
runs the statements that the repository
+ * prepares for the default table (primary key KEY).
+ */
+ private CqlSession inMemorySession(Map<String, Object[]> table) {
+ CqlSession session = mock(CqlSession.class);
+ Map<BoundStatement, Object[]> boundValues = new ConcurrentHashMap<>();
+ Map<BoundStatement, String> boundQueries = new ConcurrentHashMap<>();
+
when(session.prepare(any(SimpleStatement.class))).thenAnswer(invocation -> {
+ String query = invocation.getArgument(0,
SimpleStatement.class).getQuery().toUpperCase();
+ PreparedStatement prepared = mock(PreparedStatement.class);
+ when(prepared.bind(any(Object[].class))).thenAnswer(bind -> {
+ BoundStatement bound = mock(BoundStatement.class);
+ boundValues.put(bound, bind.getArguments());
+ boundQueries.put(bound, query);
+ return bound;
+ });
+ return prepared;
+ });
+ when(session.execute(any(BoundStatement.class))).thenAnswer(invocation
-> {
+ BoundStatement bound = invocation.getArgument(0);
+ return execute(table, boundQueries.remove(bound),
boundValues.remove(bound));
+ });
+ return session;
+ }
+
+ private static ResultSet execute(Map<String, Object[]> table, String
query, Object[] values) {
+ List<Row> rows = new ArrayList<>();
+ if (query.startsWith("INSERT")) {
+ table.put((String) values[0], values.clone());
+ } else if (query.startsWith("DELETE") && query.contains(" IF ")) {
+ table.computeIfPresent((String) values[0], (key, row) ->
values[1].equals(row[1]) ? null : row);
+ } else if (query.startsWith("DELETE")) {
+ table.remove((String) values[0]);
+ } else if (query.startsWith("SELECT") && query.contains("WHERE")) {
+ Object[] row = table.get((String) values[0]);
+ if (row != null) {
+ rows.add(row(row));
+ }
+ } else if (query.startsWith("SELECT")) {
+ table.values().forEach(row -> rows.add(row(row)));
+ } else {
+ throw new IllegalArgumentException("Unexpected statement " +
query);
+ }
+ ResultSet resultSet = mock(ResultSet.class);
+ when(resultSet.one()).thenReturn(rows.isEmpty() ? null : rows.get(0));
+ when(resultSet.all()).thenReturn(rows);
+ return resultSet;
+ }
+
+ private static Row row(Object[] values) {
+ Row row = mock(Row.class);
+ when(row.getString(anyString())).thenAnswer(invocation -> switch
(invocation.getArgument(0, String.class)) {
+ case "KEY" -> values[0];
+ case "EXCHANGE_ID" -> values[1];
+ default -> throw new
IllegalArgumentException(invocation.getArgument(0, String.class));
+ });
+ when(row.getByteBuffer("EXCHANGE")).thenAnswer(invocation ->
((ByteBuffer) values[2]).duplicate());
+ return row;
+ }
+}
diff --git
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
index 32273139cc7a..a6e4106974d8 100644
---
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
+++
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
@@ -145,11 +145,16 @@ public class NamedCassandraAggregationRepositoryIT
extends BaseCassandra {
aggregationRepository.add(context, key, exchange);
assertTrue(exists(key));
}
+ Exchange completed = new DefaultExchange(context);
+ completed.setExchangeId("Exchange_2");
+ aggregationRepository.remove(context, "Confirm_2", completed);
+ assertFalse(exists("Confirm_2"));
+ assertTrue(exists("camel-recovery:Exchange_2"));
// When
aggregationRepository.confirm(context, "Exchange_2");
// Then
+ assertFalse(exists("camel-recovery:Exchange_2"));
assertTrue(exists("Confirm_1"));
- assertFalse(exists("Confirm_2"));
assertTrue(exists("Confirm_3"));
}
@@ -172,6 +177,12 @@ public class NamedCassandraAggregationRepositoryIT extends
BaseCassandra {
}
}
+ private void removeExchange(String key) {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.setExchangeId("Exchange-" + key);
+ aggregationRepository.remove(context, key, exchange);
+ }
+
private void addExchanges(String... keys) {
for (String key : keys) {
Exchange exchange = new DefaultExchange(context);
@@ -185,12 +196,16 @@ public class NamedCassandraAggregationRepositoryIT
extends BaseCassandra {
// Given
String[] keys = { "Scan1", "Scan2" };
addExchanges(keys);
+ // aggregations in progress are not to be recovered
+
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan1"));
+
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan2"));
// When
+ removeExchange("Scan2");
Set<String> exchangeIdSet = aggregationRepository.scan(context);
// Then
- for (String key : keys) {
- assertTrue(exchangeIdSet.contains("Exchange-" + key));
- }
+ assertTrue(exchangeIdSet.contains("Exchange-Scan2"));
+ assertFalse(exchangeIdSet.contains("Exchange-Scan1"));
+
assertFalse(aggregationRepository.getKeys().contains("camel-recovery:Exchange-Scan2"));
}
@Test
@@ -198,11 +213,15 @@ public class NamedCassandraAggregationRepositoryIT
extends BaseCassandra {
// Given
String[] keys = { "Recover1", "Recover2" };
addExchanges(keys);
+ removeExchange("Recover2");
// When
+ Exchange exchange1 = aggregationRepository.recover(context,
"Exchange-Recover1");
Exchange exchange2 = aggregationRepository.recover(context,
"Exchange-Recover2");
Exchange exchange3 = aggregationRepository.recover(context,
"Exchange-Recover3");
// Then
+ assertNull(exchange1);
assertNotNull(exchange2);
+ assertEquals("Exchange-Recover2", exchange2.getExchangeId());
assertNull(exchange3);
}
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 7fbdd1071cb2..7b807df8f8a2 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
@@ -3048,6 +3048,24 @@ A poll in which the download of every file failed now
counts as an idle poll, th
could be acquired. This affects `backoffIdleThreshold`,
`sendEmptyMessageWhenIdle` (an empty message is sent) and
`greedy` (the consumer does not poll again immediately).
+=== camel-cassandraql - the aggregation repository keeps completed exchanges
for recovery
+
+`CassandraAggregationRepository` (and `NamedCassandraAggregationRepository`)
implement
+`RecoverableAggregationRepository` with recovery enabled by default, but had
no recovery store: `remove` deleted
+the completed exchange, `confirm` deleted the aggregation in progress that had
the confirmed exchange id, and `scan`
+returned the exchange ids of the aggregations still in progress. The recovery
task therefore sent aggregations
+that were still in progress (marked `CamelRedelivered`) and then deleted them,
while exchanges that failed after
+completion could never be recovered.
+
+A completed exchange is now kept in the same table, under the aggregation key
`camel-recovery:<exchange id>`, until
+it is confirmed, which is what `scan` reports and `recover` reads. `getKeys`
reports only the aggregations in
+progress. No change of the table is needed.
+
+Routes that set `useRecovery=false` are unaffected. Routes that left recovery
enabled, which is the default, no
+longer see aggregations in progress sent by the recovery task, and a completed
exchange whose processing failed is
+now recovered. The table now also holds one row per completed and not yet
confirmed exchange; those rows are deleted
+on confirmation. A time to live (`ttl`) applies to these rows as well.
+
=== camel-infinispan - the aggregation repository keeps completed exchanges
for recovery
`InfinispanAggregationRepository` implements
`RecoverableAggregationRepository`, but it had no recovery