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 f78a93999a24 CAMEL-25124: camel-core - Aggregate EIP: keep spooled
stream cache bodies until the aggregated exchange is done (#27033)
f78a93999a24 is described below
commit f78a93999a24d1fb0f1e85474b2fdc8aa9b71a9d
Author: allthingssecurity <[email protected]>
AuthorDate: Tue Sep 29 17:25:35 2026 +0530
CAMEL-25124: camel-core - Aggregate EIP: keep spooled stream cache bodies
until the aggregated exchange is done (#27033)
* CAMEL-25124: camel-core - Aggregate EIP: keep spooled stream cache bodies
until the aggregated exchange is done
* CAMEL-25124: camel-core - Aggregate EIP: release only the spool
references the aggregator took
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../TarAggregationStrategyRepositoryTest.java | 147 +++++++
.../ZipAggregationStrategyRepositoryTest.java | 147 +++++++
.../processor/aggregate/AggregateProcessor.java | 159 ++++++-
.../AggregateStreamCachingSpoolTest.java | 458 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 16 +
5 files changed, 925 insertions(+), 2 deletions(-)
diff --git
a/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
b/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
new file mode 100644
index 000000000000..f8b8c85d8954
--- /dev/null
+++
b/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
@@ -0,0 +1,147 @@
+/*
+ * 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.tarfile;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+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.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.commons.compress.archivers.tar.TarArchiveEntry;
+import org.apache.commons.compress.archivers.tar.TarArchiveInputStream;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.apache.camel.test.junit6.TestSupport.deleteDirectory;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+/**
+ * The tar file that the strategy appends to must not be deleted before the
group completes, when the aggregation
+ * repository does not keep the exchange instances (optimistic locking, or a
persistent repository).
+ */
+public class TarAggregationStrategyRepositoryTest extends CamelTestSupport {
+
+ private static final String TEST_DIR =
"target/out_TarAggregationStrategyRepositoryTest";
+
+ @BeforeEach
+ public void deleteTestDirs() {
+ deleteDirectory(TEST_DIR);
+ }
+
+ @Test
+ public void testOptimisticLocking() throws Exception {
+ sendAndAssertZip("optimistic");
+ }
+
+ @Test
+ public void testRepositoryStoringCopies() throws Exception {
+ sendAndAssertZip("copies");
+ }
+
+ private void sendAndAssertZip(String name) throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:" + name);
+ mock.expectedMessageCount(1);
+
+ template.sendBody("direct:" + name, "Hello");
+ template.sendBody("direct:" + name, "Hello again");
+ template.sendBody("direct:" + name, "Bye");
+
+ MockEndpoint.assertIsSatisfied(context);
+
+ File[] files = new File(TEST_DIR, name).listFiles();
+ assertNotNull(files);
+ assertEquals(1, files.length);
+ int count = 0;
+ try (TarArchiveInputStream tin = new TarArchiveInputStream(new
FileInputStream(files[0]))) {
+ for (TarArchiveEntry te = tin.getNextEntry(); te != null; te =
tin.getNextEntry()) {
+ count++;
+ }
+ }
+ assertEquals(3, count, "Tar file should contain 3 files");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:optimistic")
+ .aggregate(constant(true), strategy())
+ .aggregationRepository(new
MemoryAggregationRepository(true)).optimisticLocking()
+ .completionSize(3)
+ .to("file:" + TEST_DIR + "/optimistic")
+ .to("mock:optimistic");
+
+ from("direct:copies")
+ .aggregate(constant(true), strategy())
+ .aggregationRepository(new CopyingRepository())
+ .completionSize(3)
+ .to("file:" + TEST_DIR + "/copies")
+ .to("mock:copies");
+ }
+ };
+ }
+
+ private static TarAggregationStrategy strategy() {
+ TarAggregationStrategy strategy = new TarAggregationStrategy();
+ strategy.setParentDir(TEST_DIR + "/temp");
+ return strategy;
+ }
+
+ /**
+ * A repository that stores a copy of the exchange, as the persistent
repositories do.
+ */
+ private static final class CopyingRepository extends ServiceSupport
implements AggregationRepository {
+ private final Map<String, Exchange> exchanges = new
ConcurrentHashMap<>();
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
exchange) {
+ return exchanges.put(key, exchange.copy());
+ }
+
+ @Override
+ public Exchange get(CamelContext camelContext, String key) {
+ Exchange exchange = exchanges.get(key);
+ return exchange != null ? exchange.copy() : null;
+ }
+
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ exchanges.remove(key);
+ }
+
+ @Override
+ public void confirm(CamelContext camelContext, String exchangeId) {
+ // noop
+ }
+
+ @Override
+ public Set<String> getKeys() {
+ return exchanges.keySet();
+ }
+ }
+}
diff --git
a/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
b/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
new file mode 100644
index 000000000000..bf28b2944cab
--- /dev/null
+++
b/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
@@ -0,0 +1,147 @@
+/*
+ * 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.zipfile;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.zip.ZipEntry;
+import java.util.zip.ZipInputStream;
+
+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.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.apache.camel.test.junit6.TestSupport.deleteDirectory;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+/**
+ * The zip file that the strategy appends to must not be deleted before the
group completes, when the aggregation
+ * repository does not keep the exchange instances (optimistic locking, or a
persistent repository).
+ */
+public class ZipAggregationStrategyRepositoryTest extends CamelTestSupport {
+
+ private static final String TEST_DIR =
"target/out_ZipAggregationStrategyRepositoryTest";
+
+ @BeforeEach
+ public void deleteTestDirs() {
+ deleteDirectory(TEST_DIR);
+ }
+
+ @Test
+ public void testOptimisticLocking() throws Exception {
+ sendAndAssertZip("optimistic");
+ }
+
+ @Test
+ public void testRepositoryStoringCopies() throws Exception {
+ sendAndAssertZip("copies");
+ }
+
+ private void sendAndAssertZip(String name) throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:" + name);
+ mock.expectedMessageCount(1);
+
+ template.sendBody("direct:" + name, "Hello");
+ template.sendBody("direct:" + name, "Hello again");
+ template.sendBody("direct:" + name, "Bye");
+
+ MockEndpoint.assertIsSatisfied(context);
+
+ File[] files = new File(TEST_DIR, name).listFiles();
+ assertNotNull(files);
+ assertEquals(1, files.length);
+ int count = 0;
+ try (ZipInputStream zin = new ZipInputStream(new
FileInputStream(files[0]))) {
+ for (ZipEntry ze = zin.getNextEntry(); ze != null; ze =
zin.getNextEntry()) {
+ count++;
+ }
+ }
+ assertEquals(3, count, "Zip file should contain 3 files");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:optimistic")
+ .aggregate(constant(true), strategy())
+ .aggregationRepository(new
MemoryAggregationRepository(true)).optimisticLocking()
+ .completionSize(3)
+ .to("file:" + TEST_DIR + "/optimistic")
+ .to("mock:optimistic");
+
+ from("direct:copies")
+ .aggregate(constant(true), strategy())
+ .aggregationRepository(new CopyingRepository())
+ .completionSize(3)
+ .to("file:" + TEST_DIR + "/copies")
+ .to("mock:copies");
+ }
+ };
+ }
+
+ private static ZipAggregationStrategy strategy() {
+ ZipAggregationStrategy strategy = new ZipAggregationStrategy();
+ strategy.setParentDir(TEST_DIR + "/temp");
+ return strategy;
+ }
+
+ /**
+ * A repository that stores a copy of the exchange, as the persistent
repositories do.
+ */
+ private static final class CopyingRepository extends ServiceSupport
implements AggregationRepository {
+ private final Map<String, Exchange> exchanges = new
ConcurrentHashMap<>();
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
exchange) {
+ return exchanges.put(key, exchange.copy());
+ }
+
+ @Override
+ public Exchange get(CamelContext camelContext, String key) {
+ Exchange exchange = exchanges.get(key);
+ return exchange != null ? exchange.copy() : null;
+ }
+
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ exchanges.remove(key);
+ }
+
+ @Override
+ public void confirm(CamelContext camelContext, String exchangeId) {
+ // noop
+ }
+
+ @Override
+ public Set<String> getKeys() {
+ return exchanges.keySet();
+ }
+ }
+}
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 574c03376e9b..93392d6c66a3 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
@@ -16,7 +16,9 @@
*/
package org.apache.camel.processor.aggregate;
+import java.io.IOException;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -48,6 +50,7 @@ import org.apache.camel.Predicate;
import org.apache.camel.Processor;
import org.apache.camel.ProducerTemplate;
import org.apache.camel.ShutdownRunningTask;
+import org.apache.camel.StreamCache;
import org.apache.camel.TimeoutMap;
import org.apache.camel.Traceable;
import org.apache.camel.processor.BaseProcessorSupport;
@@ -62,12 +65,15 @@ import org.apache.camel.spi.RouteIdAware;
import org.apache.camel.spi.ShutdownAware;
import org.apache.camel.spi.StepIdAware;
import org.apache.camel.spi.Synchronization;
+import org.apache.camel.support.DefaultExchange;
import org.apache.camel.support.DefaultTimeoutMap;
import org.apache.camel.support.ExchangeHelper;
import org.apache.camel.support.KeyValueAggregationRepository;
import org.apache.camel.support.LRUCacheFactory;
import org.apache.camel.support.LoggingExceptionHandler;
import org.apache.camel.support.NoLock;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.support.UnitOfWorkHelper;
import org.apache.camel.support.service.ServiceHelper;
import org.apache.camel.util.ObjectHelper;
import org.apache.camel.util.StopWatch;
@@ -135,6 +141,9 @@ public class AggregateProcessor extends BaseProcessorSupport
// aggregated exchanges being completed (exchange id -> count): moved to
the completed store of a recoverable
// repository, but not yet registered in inProgressCompleteExchanges by
onSubmitCompletion
private final Map<String, Integer> completingExchanges = new
ConcurrentHashMap<>();
+ // the references to spooled stream caches that the groups in the
repository hold, by correlation key (only when
+ // isKeepingReferences, and guarded by the lock)
+ private final Map<String, SpooledStreamCaches> groupSpooledStreamCaches =
new HashMap<>();
private final Map<String, RedeliveryData> redeliveryState = new
ConcurrentHashMap<>();
private final AggregateProcessorStatistics statistics = new Statistics();
@@ -451,6 +460,31 @@ public class AggregateProcessor extends
BaseProcessorSupport
removeFlagCompleteAllGroups(copy);
removeFlagCompleteAllGroupsInclusive(copy);
+ // a stream cache spooled to disk is deleted when the incoming
exchange is done, but the group keeps the body
+ // for later, so the copy takes its own reference, which the
aggregator releases when it is done with the body
+ SpooledStreamCaches spooled = null;
+ if (copy.getIn().getBody() instanceof StreamCache sc &&
!sc.inMemory()) {
+ // the copy is independent of the unit of work that the incoming
exchange may release its stream caches with
+ // (such as the parent of a split), as the wire tap does
+ copy.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK);
+ try {
+ StreamCache copied = sc.copy(copy);
+ if (copied != null) {
+ copy.getIn().setBody(copied);
+ }
+ } catch (IOException e) {
+ UnitOfWorkHelper.doneSynchronizations(copy,
copy.getExchangeExtension().handoverCompletions());
+ exchange.setException(e);
+ callback.done(sync);
+ return sync;
+ }
+ // the copy has no on completions of its own, so these are only
the releases of the reference, which are
+ // kept apart from the on completions that the aggregation
strategy may add to the exchanges (such as the
+ // ZipAggregationStrategy deleting its zip file when the
aggregated exchange is done)
+ spooled = new SpooledStreamCaches();
+ spooled.add(copy.getExchangeExtension().handoverCompletions());
+ }
+
List<Exchange> aggregated = new ArrayList<>();
lock.lock();
try {
@@ -459,11 +493,17 @@ public class AggregateProcessor extends
BaseProcessorSupport
if (closedCorrelationKeys != null &&
closedCorrelationKeys.containsKey(key)) {
throw new ClosedCorrelationKeyException(key, exchange);
}
- doAggregation(key, copy, aggregated);
+ doAggregation(key, copy, spooled, aggregated);
} catch (CamelExchangeException e) {
exchange.setException(e);
} finally {
lock.unlock();
+ if (spooled != null) {
+ // release the reference unless it was handed over to the
group or the aggregated exchange: the copy
+ // failed, is retried with a new copy due to optimistic
locking, or is discarded, or the repository has
+ // read the body when the copy was added (see
isKeepingReferences)
+ spooled.onDone(copy);
+ }
// we are completed so submit to completion outside the lock. This
must also be done when the aggregation
// failed, or must be retried due to optimistic locking, after a
group was completed (such as a group
// completed by pre-completion), as that group has already been
removed from the repository
@@ -524,10 +564,13 @@ public class AggregateProcessor extends
BaseProcessorSupport
*
* @param key the correlation key
* @param newExchange the exchange
+ * @param spooled the references that the
exchange holds to spooled stream caches,
+ * which are taken when
they are handed over, or <tt>null</tt>
* @param list the list to add the
aggregated exchange(s) which are complete to
* @throws org.apache.camel.CamelExchangeException is thrown if error
aggregating
*/
- private void doAggregation(String key, Exchange newExchange,
List<Exchange> list) throws CamelExchangeException {
+ private void doAggregation(String key, Exchange newExchange,
SpooledStreamCaches spooled, List<Exchange> list)
+ throws CamelExchangeException {
LOG.trace("onAggregation +++ start +++ with correlation key: {}", key);
String complete = null;
@@ -667,9 +710,19 @@ public class AggregateProcessor extends
BaseProcessorSupport
}
// only need to update aggregation repository if we are not
complete
doAggregationRepositoryAdd(newExchange.getContext(), key,
originalExchange, answer);
+ if (spooled != null && isKeepingReferences()) {
+ // the group keeps the references until it is completed (see
doOnCompletion)
+ groupSpooledStreamCaches.computeIfAbsent(key, k -> new
SpooledStreamCaches()).add(spooled.take());
+ }
} else {
// if we are complete then add the answer to the list
doAggregationComplete(complete, list, key, originalExchange,
answer, aggregateFailed);
+ if (spooled != null && containsInstance(list, answer)) {
+ // the aggregated exchange takes over the references
(onCompletion has handed over those of the group)
+ SpooledStreamCaches release = new SpooledStreamCaches();
+ release.add(spooled.take());
+ answer.getExchangeExtension().addOnCompletion(release);
+ }
}
LOG.trace("onAggregation +++ end +++ with correlation key: {}", key);
@@ -711,6 +764,18 @@ public class AggregateProcessor extends
BaseProcessorSupport
}
}
+ /**
+ * Whether the groups in the repository keep the references to the spooled
stream caches of their exchanges until
+ * they are completed. This is only the case with the memory repository,
which keeps the exchange instances (and
+ * their bodies), and without optimistic locking, as all access to the
groups is then under the lock. A persistent
+ * repository has read the body when the exchange is added, so the
reference is released after the add. With
+ * optimistic locking the reference is released after the add as well, and
only the body of the exchange that
+ * completes the group is kept.
+ */
+ private boolean isKeepingReferences() {
+ return !optimisticLocking && aggregationRepository instanceof
MemoryAggregationRepository;
+ }
+
protected void doAggregationRepositoryAdd(
CamelContext camelContext, String key, Exchange oldExchange,
Exchange newExchange) {
LOG.trace("In progress aggregated oldExchange: {}, newExchange: {}
with correlation key: {}", oldExchange, newExchange,
@@ -894,6 +959,11 @@ public class AggregateProcessor extends
BaseProcessorSupport
if (original != null) {
// remove from repository as its completed, we do this first as to
trigger any OptimisticLockingException's
aggregationRepository.remove(aggregated.getContext(), key,
original);
+ // the aggregated exchange takes over the references that the
group holds to spooled stream caches
+ SpooledStreamCaches spooled = takeGroupSpooledStreamCaches(key);
+ if (spooled != null) {
+ aggregated.getExchangeExtension().addOnCompletion(spooled);
+ }
}
// cleanup timeout map if it was a incoming exchange which triggered
the timeout (and not the timeout checker)
@@ -949,9 +1019,82 @@ public class AggregateProcessor extends
BaseProcessorSupport
aggregationRepository.confirm(aggregated.getContext(),
aggregated.getExchangeId());
// and remove redelivery state as well
redeliveryState.remove(aggregated.getExchangeId());
+ // and release the references that the discarded exchange holds to
spooled stream caches (only those, as the
+ // other on completions of the exchange are left as they are when it
is not discarded)
+ releaseSpooledStreamCaches(aggregated);
// the completion was from timeout and we should just discard it
}
+ /**
+ * Takes the references that the group holds to spooled stream caches,
when the group is removed from the
+ * repository.
+ */
+ private SpooledStreamCaches takeGroupSpooledStreamCaches(String key) {
+ return isKeepingReferences() ? groupSpooledStreamCaches.remove(key) :
null;
+ }
+
+ /**
+ * Releases the references that the aggregator handed over to a discarded
exchange (see doOnCompletion).
+ */
+ private static void releaseSpooledStreamCaches(Exchange exchange) {
+ List<Synchronization> completions =
exchange.getExchangeExtension().handoverCompletions();
+ if (completions != null) {
+ for (Synchronization completion : completions) {
+ if (completion instanceof SpooledStreamCaches spooled) {
+ spooled.onDone(exchange);
+ } else {
+
exchange.getExchangeExtension().addOnCompletion(completion);
+ }
+ }
+ }
+ }
+
+ /**
+ * The on completions that release references that the aggregator took to
spooled stream caches, so the bodies can
+ * be read until the aggregator is done with them. They are kept apart
from the other on completions of the
+ * exchanges, which the aggregator leaves as they are. When added to an
exchange as on completion, they are released
+ * when the exchange is done.
+ */
+ private static final class SpooledStreamCaches extends
SynchronizationAdapter {
+ private List<Synchronization> releases;
+
+ void add(List<Synchronization> synchronizations) {
+ if (synchronizations != null && !synchronizations.isEmpty()) {
+ if (releases == null) {
+ releases = new ArrayList<>(synchronizations);
+ } else {
+ releases.addAll(synchronizations);
+ }
+ }
+ }
+
+ List<Synchronization> take() {
+ List<Synchronization> answer = releases;
+ releases = null;
+ return answer;
+ }
+
+ @Override
+ public void onDone(Exchange exchange) {
+ // release only once
+ UnitOfWorkHelper.doneSynchronizations(exchange, take());
+ }
+
+ @Override
+ public String toString() {
+ return "AggregateOnCompletion[SpooledStreamCaches]";
+ }
+ }
+
+ private static boolean containsInstance(List<Exchange> list, Exchange
exchange) {
+ for (Exchange e : list) {
+ if (e == exchange) {
+ return true;
+ }
+ }
+ return false;
+ }
+
private void onSubmitCompletion(final String key, final Exchange exchange)
{
LOG.debug("Aggregation complete for correlation key {} sending
aggregated exchange: {}", key, exchange);
@@ -1902,6 +2045,18 @@ public class AggregateProcessor extends
BaseProcessorSupport
// shutdown aggregation repository and the strategy
ServiceHelper.stopAndShutdownServices(aggregationRepository,
aggregationStrategy);
+ // the memory repository has dropped the groups, so release the
references they held to spooled stream caches
+ if (!groupSpooledStreamCaches.isEmpty()) {
+ lock.lock();
+ try {
+ Exchange dummy = new DefaultExchange(camelContext);
+ groupSpooledStreamCaches.values().forEach(spooled ->
spooled.onDone(dummy));
+ groupSpooledStreamCaches.clear();
+ } finally {
+ lock.unlock();
+ }
+ }
+
// cleanup when shutting down
inProgressCompleteExchanges.clear();
completingExchanges.clear();
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
new file mode 100644
index 000000000000..dc334841518b
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
@@ -0,0 +1,458 @@
+/*
+ * 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.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.InputStream;
+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.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.CamelContext;
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.processor.aggregate.AggregateController;
+import org.apache.camel.processor.aggregate.ClosedCorrelationKeyException;
+import org.apache.camel.processor.aggregate.DefaultAggregateController;
+import org.apache.camel.processor.aggregate.GroupedBodyAggregationStrategy;
+import org.apache.camel.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.processor.aggregate.UseLatestAggregationStrategy;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.spi.OptimisticLockingAggregationRepository;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.support.service.ServiceSupport;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * The Aggregate EIP must keep the bodies that stream caching spooled to disk
until the aggregated exchange is done with
+ * them, and must not leave the spool files behind, also when a group or an
incoming exchange is discarded.
+ */
+public class AggregateStreamCachingSpoolTest extends ContextTestSupport {
+
+ private static final byte[] DATA = createData(16 * 1024);
+
+ // the aggregated exchange is processed only when the incoming exchanges
are done, so the spool files would already
+ // be deleted if the aggregator did not hold its own references to them
+ private final CountDownLatch sendersDone = new CountDownLatch(1);
+ private final AggregateController controller = new
DefaultAggregateController();
+ private final FailOnceRepository failOnce = new FailOnceRepository();
+ private final FailOnceRepository failOnceAndClose = new
FailOnceRepository();
+ // whether the on completion that the aggregation strategy added to the
first exchange of the group has run
+ private final AtomicBoolean strategyCompletionDone = new AtomicBoolean();
+
+ @Test
+ public void testGroupedBodies() throws Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(1);
+
+ template.sendBody("direct:grouped", stream());
+ template.sendBody("direct:grouped", stream());
+ template.sendBody("direct:grouped", stream());
+ sendersDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ List<?> bodies =
result.getReceivedExchanges().get(0).getMessage().getBody(List.class);
+ assertEquals(3, bodies.size());
+ for (Object body : bodies) {
+ assertArrayEquals(DATA, (byte[]) body);
+ }
+ assertSpoolDirectoryEmpty();
+ }
+
+ @Test
+ public void testCompletionTimeout() throws Exception {
+ sendAndAssertBody("direct:timeout");
+ }
+
+ @Test
+ public void testCompletionSizeOne() throws Exception {
+ // the aggregated exchange is sent on the aggregator's own thread, so
even the body of the exchange that
+ // completes the group must not depend on the incoming exchange
+ sendAndAssertBody("direct:size");
+ }
+
+ @Test
+ public void testOptimisticLockingRetry() throws Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(1);
+
+ // the first attempt to add the first exchange fails, and the retry
takes a new copy
+ template.sendBody("direct:optimistic", stream());
+ template.sendBody("direct:optimistic", stream());
+ sendersDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ assertEquals(2, failOnce.calls.get());
+ assertArrayEquals(DATA,
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+ assertSpoolDirectoryEmpty();
+ }
+
+ @Test
+ public void testClosedCorrelationKeyOnRetry() throws Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(1);
+ sendersDone.countDown();
+
+ // while the first attempt to add this exchange fails, another
exchange completes and closes the group, so the
+ // retry finds the correlation key closed
+ CamelExecutionException e = assertThrows(CamelExecutionException.class,
+ () -> template.sendBody("direct:closed", stream()));
+ assertInstanceOf(ClosedCorrelationKeyException.class, e.getCause());
+
+ assertMockEndpointsSatisfied();
+ assertArrayEquals(DATA,
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+ assertSpoolDirectoryEmpty();
+ }
+
+ @Test
+ public void testDiscardOnAggregationFailure() throws Exception {
+ getMockEndpoint("mock:result").expectedMessageCount(0);
+
+ // the first exchange fails to aggregate, and is discarded
+ template.sendBodyAndHeader("direct:failure", stream(), "fail", true);
+ // the second exchange starts a group, and the third fails to
aggregate, which discards the group
+ template.sendBody("direct:failure", stream());
+ template.sendBodyAndHeader("direct:failure", stream(), "fail", true);
+ sendersDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ assertSpoolDirectoryEmpty();
+ }
+
+ @Test
+ public void testDiscardOnCompletionTimeout() throws Exception {
+ getMockEndpoint("mock:result").expectedMessageCount(0);
+
+ template.sendBody("direct:discardTimeout", stream());
+ sendersDone.countDown();
+
+ assertSpoolDirectoryEmpty();
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testForceDiscardingOfGroup() throws Exception {
+ getMockEndpoint("mock:result").expectedMessageCount(0);
+
+ template.sendBody("direct:forceDiscard", stream());
+ template.sendBody("direct:forceDiscard", stream());
+ sendersDone.countDown();
+ assertEquals(1, controller.forceDiscardingOfGroup("group"));
+
+ assertMockEndpointsSatisfied();
+ assertSpoolDirectoryEmpty();
+ }
+
+ @Test
+ public void testRepositoryStoringCopies() throws Exception {
+ // like the persistent repositories, the repository reads the body
when the exchange is added
+ sendAndAssertBody("direct:copies");
+ }
+
+ @Test
+ public void testMemoryRepositoryStoringCopies() throws Exception {
+ // a subclass of the memory repository that does not keep the exchange
instance
+ sendAndAssertBody("direct:memoryCopies");
+ }
+
+ @Test
+ public void testStrategyCompletionWithOptimisticLocking() throws Exception
{
+ sendAndAssertStrategyCompletion("direct:strategyCompletionOptimistic");
+ }
+
+ @Test
+ public void testStrategyCompletionWithRepositoryStoringCopies() throws
Exception {
+ sendAndAssertStrategyCompletion("direct:strategyCompletionCopies");
+ }
+
+ @Test
+ public void testStrategyCompletionWithMemoryRepository() throws Exception {
+ sendAndAssertStrategyCompletion("direct:strategyCompletionMemory");
+ }
+
+ private void sendAndAssertStrategyCompletion(String uri) throws Exception {
+ // the aggregator must only release its own references to the spooled
stream caches, and leave the on
+ // completions that the aggregation strategy adds to the exchanges
alone (such as the ZipAggregationStrategy
+ // that deletes its zip file when the aggregated exchange is done)
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(1);
+
+ template.sendBody(uri, stream());
+ template.sendBody(uri, stream());
+ template.sendBodyAndHeader(uri, stream(), "last", true);
+ sendersDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ assertArrayEquals(DATA,
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+ assertSpoolDirectoryEmpty();
+ if (uri.endsWith("Memory")) {
+ // the memory repository keeps the exchange, so the on completion
runs when the aggregated exchange is done
+ Awaitility.await().atMost(5,
TimeUnit.SECONDS).untilTrue(strategyCompletionDone);
+ }
+ }
+
+ private void sendAndAssertBody(String uri) throws Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(1);
+
+ template.sendBody(uri, stream());
+ sendersDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ assertArrayEquals(DATA,
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+ assertSpoolDirectoryEmpty();
+ }
+
+ private void assertSpoolDirectoryEmpty() {
+ File spoolDir = testDirectory().toFile();
+ Awaitility.await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> {
+ String[] files = spoolDir.list();
+ assertNotNull(files);
+ assertEquals(0, files.length, "Spool files left behind: " +
List.of(files));
+ });
+ }
+
+ private static InputStream stream() {
+ // a stream that is not converted to an in-memory cache, so it is
spooled to disk
+ return new BufferedInputStream(new ByteArrayInputStream(DATA));
+ }
+
+ private static byte[] createData(int size) {
+ byte[] data = new byte[size];
+ for (int i = 0; i < size; i++) {
+ data[i] = (byte) ('a' + i % 26);
+ }
+ return data;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
context.getStreamCachingStrategy().setSpoolDirectory(testDirectory().toFile());
+ context.getStreamCachingStrategy().setSpoolEnabled(true);
+ context.getStreamCachingStrategy().setSpoolThreshold(1024);
+
context.getStreamCachingStrategy().setRemoveSpoolDirectoryWhenStopping(false);
+ context.setStreamCaching(true);
+
+ // reads the bodies after the incoming exchanges are done
+ Processor read = e -> {
+ sendersDone.await(20, TimeUnit.SECONDS);
+ Object body = e.getMessage().getBody();
+ if (body instanceof List<?> list) {
+ List<byte[]> answer = new ArrayList<>();
+ for (Object o : list) {
+
answer.add(e.getContext().getTypeConverter().mandatoryConvertTo(byte[].class,
e, o));
+ }
+ e.getMessage().setBody(answer);
+ } else {
+
e.getMessage().setBody(e.getMessage().getMandatoryBody(byte[].class));
+ }
+ };
+
+ // adds an on completion to the first exchange of the group,
which must not run before the group completes
+ AggregationStrategy addCompletion = (oldExchange, newExchange)
-> {
+ if (oldExchange == null) {
+ strategyCompletionDone.set(false);
+ newExchange.getExchangeExtension().addOnCompletion(new
SynchronizationAdapter() {
+ @Override
+ public void onDone(Exchange exchange) {
+ strategyCompletionDone.set(true);
+ }
+ });
+ return newExchange;
+ }
+ if (strategyCompletionDone.get()) {
+ throw new IllegalStateException("The on completion of
the aggregation strategy has already run");
+ }
+
oldExchange.getMessage().setBody(newExchange.getMessage().getBody());
+ return oldExchange;
+ };
+
+ AggregationStrategy failOnHeader = (oldExchange, newExchange)
-> {
+ if (newExchange.getMessage().getHeader("fail") != null) {
+ throw new IllegalArgumentException("Forced");
+ }
+ return newExchange;
+ };
+
+ from("direct:grouped")
+ .aggregate(constant("group"), new
GroupedBodyAggregationStrategy()).completionSize(3)
+ .process(read).to("mock:result");
+
+ from("direct:timeout")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+ .process(read).to("mock:result");
+
+ from("direct:size")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy()).completionSize(1)
+ .process(read).to("mock:result");
+
+ from("direct:optimistic")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+
.aggregationRepository(failOnce).optimisticLocking().completionSize(2)
+ .process(read).to("mock:result");
+
+ from("direct:closed")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+
.aggregationRepository(failOnceAndClose).optimisticLocking()
+
.completionPredicate(header("complete").isNotNull()).closeCorrelationKeyOnCompletion(100)
+ .process(read).to("mock:result");
+
+ from("direct:failure")
+ .aggregate(constant("group"),
failOnHeader).discardOnAggregationFailure().completionSize(5)
+ .to("mock:result");
+
+ from("direct:discardTimeout")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+
.completionTimeout(100).completionTimeoutCheckerInterval(10).discardOnCompletionTimeout()
+ .to("mock:result");
+
+ from("direct:forceDiscard")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy()).aggregateController(controller)
+ .completionSize(5)
+ .to("mock:result");
+
+ from("direct:strategyCompletionOptimistic")
+ .aggregate(constant("group"), addCompletion)
+ .aggregationRepository(new
MemoryAggregationRepository(true)).optimisticLocking()
+
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+ .process(read).to("mock:result");
+
+ from("direct:strategyCompletionCopies")
+ .aggregate(constant("group"), addCompletion)
+ .aggregationRepository(new CopyingRepository())
+
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+ .process(read).to("mock:result");
+
+ from("direct:strategyCompletionMemory")
+ .aggregate(constant("group"), addCompletion)
+
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+ .process(read).to("mock:result");
+
+ from("direct:memoryCopies")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+ .aggregationRepository(new
MemoryAggregationRepository() {
+ @Override
+ public Exchange add(CamelContext camelContext,
String key, Exchange exchange) {
+ Exchange copy = exchange.copy();
+
copy.getMessage().setBody(exchange.getMessage().getBody(byte[].class));
+ return super.add(camelContext, key, copy);
+ }
+ })
+
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+ .process(read).to("mock:result");
+
+ from("direct:copies")
+ .aggregate(constant("group"), new
UseLatestAggregationStrategy())
+ .aggregationRepository(new CopyingRepository())
+
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+ .process(read).to("mock:result");
+ }
+ };
+ }
+
+ /**
+ * An optimistic locking repository whose first add fails. When it closes
the group, it first sends an exchange that
+ * completes the group, which closes the correlation key.
+ */
+ private final class FailOnceRepository extends MemoryAggregationRepository
{
+ private final AtomicInteger calls = new AtomicInteger();
+
+ private FailOnceRepository() {
+ super(true);
+ }
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
oldExchange, Exchange newExchange) {
+ if (calls.incrementAndGet() == 1) {
+ if (this == failOnceAndClose) {
+ template.sendBodyAndHeader("direct:closed", stream(),
"complete", true);
+ }
+ throw new
OptimisticLockingAggregationRepository.OptimisticLockingException();
+ }
+ return super.add(camelContext, key, oldExchange, newExchange);
+ }
+ }
+
+ /**
+ * A repository that stores a copy of the exchange with the body read into
memory, as the persistent repositories
+ * store a serialized copy, and returns a new exchange instance on every
get.
+ */
+ private static final class CopyingRepository extends ServiceSupport
implements AggregationRepository {
+ private final Map<String, Exchange> exchanges = new
ConcurrentHashMap<>();
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
exchange) {
+ Exchange old = exchanges.put(key, copy(exchange,
exchange.getMessage().getBody(byte[].class)));
+ return old != null ? copy(old, old.getMessage().getBody()) : null;
+ }
+
+ @Override
+ public Exchange get(CamelContext camelContext, String key) {
+ Exchange exchange = exchanges.get(key);
+ return exchange != null ? copy(exchange,
exchange.getMessage().getBody()) : null;
+ }
+
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ exchanges.remove(key);
+ }
+
+ @Override
+ public void confirm(CamelContext camelContext, String exchangeId) {
+ // noop
+ }
+
+ @Override
+ public Set<String> getKeys() {
+ return exchanges.keySet();
+ }
+
+ private static Exchange copy(Exchange exchange, Object body) {
+ Exchange copy = new DefaultExchange(exchange.getContext());
+ copy.setExchangeId(exchange.getExchangeId());
+ copy.getProperties().putAll(exchange.getProperties());
+ copy.getMessage().setHeaders(exchange.getMessage().getHeaders());
+ copy.getMessage().setBody(body);
+ return copy;
+ }
+ }
+}
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 7e4a2cd8edcb..89c1cda4e272 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
@@ -211,6 +211,22 @@ not complete by timeout until another exchange that has a
completion timeout arr
`AggregationRepository` that does not keep exchange properties never completes
a group by timeout. This only affects
optimistic locking combined with `completionTimeoutExpression`.
+=== Aggregate EIP - stream caching with spooling to disk
+
+When stream caching spools message bodies to disk, the aggregator now keeps
the spool files of the bodies it
+aggregates until the aggregated exchange is done. Prior to Camel 4.23 a spool
file was deleted as soon as the incoming
+exchange was done, so the aggregated exchange could fail to read the bodies
with a `NoSuchFileException`.
+
+With the default `MemoryAggregationRepository` (without optimistic locking)
the spool files of all exchanges of a group
+are now kept until the group completes and the aggregated exchange is done, so
more disk space can be used in the spool
+directory than before. For example `UseLatestAggregationStrategy` with a large
`completionSize` now keeps a spool file
+for every exchange of the group, although only the last body is used. A group
that never completes keeps its spool
+files until the aggregator is shut down, such as when the route is removed or
`CamelContext` is stopped.
+
+With optimistic locking, or a persistent aggregation repository, the spool
file of an exchange is released when the
+exchange has been added to the repository (a persistent repository has read
the body at that time), and only the
+spool file of the exchange that completes a group is kept until the aggregated
exchange is done.
+
=== Variable Receive
When an EIP with `variableReceive` stores a message into a variable that
already holds a message, the header variables