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 be76af8efb2c CAMEL-24990: camel-core - Multicast, Split and Recipient
List EIPs: fix bugs found in a deep review (#26846)
be76af8efb2c is described below
commit be76af8efb2ca4ac0bc936e8dedc5c3d3ce6d97a
Author: Claus Ibsen <[email protected]>
AuthorDate: Sun Sep 27 10:09:47 2026 +0200
CAMEL-24990: camel-core - Multicast, Split and Recipient List EIPs: fix
bugs found in a deep review (#26846)
- A streaming split that finished early (stopOnException or timeout)
read the rest of the input in doDone to release the exchanges, and
ran onPrepare on every remaining part. Only a collection of pairs is
released now.
- A transacted recipient list shared one transaction context data map
across all transactions. It is now created per exchange.
- When a recipient could not be resolved, the producers already
acquired for the recipients before it were never released.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../apache/camel/processor/MulticastProcessor.java | 4 +-
.../camel/processor/RecipientListProcessor.java | 90 ++++++++--------
.../java/org/apache/camel/processor/Splitter.java | 7 +-
.../RecipientListInvalidEndpointReleaseTest.java | 117 +++++++++++++++++++++
.../RecipientListTransactedContextDataTest.java | 81 ++++++++++++++
...itterStreamingStopOnExceptionReadAheadTest.java | 78 ++++++++++++++
6 files changed, 327 insertions(+), 50 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/MulticastProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/MulticastProcessor.java
index b42093241549..243129e08ac9 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/MulticastProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/MulticastProcessor.java
@@ -981,8 +981,10 @@ public class MulticastProcessor extends
BaseProcessorSupport
}
}
- if (processorExchangeFactory != null && pairs != null) {
+ if (processorExchangeFactory != null && pairs instanceof Collection) {
// the exchanges on the pairs was created with a factory, so they
should be released
+ // (only when the pairs are a collection, as iterating a streaming
iterable would read, prepare and create
+ // exchanges for all the remaining parts, such as after
stopOnException or a timeout)
try {
for (ProcessorExchangePair pair : pairs) {
processorExchangeFactory.release(pair.getExchange());
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/RecipientListProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/RecipientListProcessor.java
index 664dce0c6085..8f3a65b31cfd 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/RecipientListProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/RecipientListProcessor.java
@@ -16,7 +16,6 @@
*/
package org.apache.camel.processor;
-import java.lang.reflect.Array;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Iterator;
@@ -74,7 +73,6 @@ public class RecipientListProcessor extends
MulticastProcessor {
private final String delimiter;
private final ProducerCache producerCache;
private int cacheSize;
- private Map<String, Object> txData;
/**
* Class that represent each step in the recipient list to do
@@ -253,52 +251,50 @@ public class RecipientListProcessor extends
MulticastProcessor {
recipientList = expression.evaluate(exchange, Object.class);
}
- // optimize for recipient without need for using delimiter
- // (if its list/collection/array type)
- if (recipientList instanceof List<?> col) {
- int size = col.size();
- List<ProcessorExchangePair> result = new ArrayList<>(size);
- int index = 0;
- for (Object recipient : col) {
- index = doCreateProcessorExchangePairs(exchange, recipient,
result, index);
- }
- return result;
- } else if (recipientList instanceof Collection<?> col) {
- int size = col.size();
- List<ProcessorExchangePair> result = new ArrayList<>(size);
+ // each exchange (transaction) has its own transaction context data,
shared by its copies
+ Map<String, Object> txData = exchange.isTransacted() ? new
ConcurrentHashMap<>() : null;
+
+ List<ProcessorExchangePair> result = new ArrayList<>();
+ try {
int index = 0;
- for (Object recipient : col) {
- index = doCreateProcessorExchangePairs(exchange, recipient,
result, index);
+ // optimize for recipient without need for using delimiter
+ // (if its collection/array type)
+ if (recipientList instanceof Collection<?> col) {
+ for (Object recipient : col) {
+ index = doCreateProcessorExchangePairs(exchange,
recipient, result, index, txData);
+ }
+ } else if (recipientList != null &&
recipientList.getClass().isArray()) {
+ for (Object recipient : (Object[]) recipientList) {
+ index = doCreateProcessorExchangePairs(exchange,
recipient, result, index, txData);
+ }
+ } else {
+ // okay we have to use iterator based separated by delimiter
+ Iterator<?> iter;
+ if (delimiter != null &&
delimiter.equalsIgnoreCase(IGNORE_DELIMITER_MARKER)) {
+ iter = ObjectHelper.createIterator(recipientList, null);
+ } else {
+ iter = ObjectHelper.createIterator(recipientList,
delimiter);
+ }
+ while (iter.hasNext()) {
+ index = doCreateProcessorExchangePairs(exchange,
iter.next(), result, index, txData);
+ }
}
- return result;
- } else if (recipientList != null &&
recipientList.getClass().isArray()) {
- Object[] arr = (Object[]) recipientList;
- int size = Array.getLength(recipientList);
- List<ProcessorExchangePair> result = new ArrayList<>(size);
- int index = 0;
- for (Object recipient : arr) {
- index = doCreateProcessorExchangePairs(exchange, recipient,
result, index);
+ } catch (Exception e) {
+ // a recipient could not be resolved, so release the producers
acquired for the recipients before it,
+ // as the recipient list is not sent to any of them
+ for (ProcessorExchangePair pair : result) {
+ if (pair instanceof RecipientProcessorExchangePair rpair) {
+ rpair.releaseIfNotBegun();
+ }
}
- return result;
- }
-
- // okay we have to use iterator based separated by delimiter
- Iterator<?> iter;
- if (delimiter != null &&
delimiter.equalsIgnoreCase(IGNORE_DELIMITER_MARKER)) {
- iter = ObjectHelper.createIterator(recipientList, null);
- } else {
- iter = ObjectHelper.createIterator(recipientList, delimiter);
- }
- List<ProcessorExchangePair> result = new ArrayList<>();
- int index = 0;
- while (iter.hasNext()) {
- index = doCreateProcessorExchangePairs(exchange, iter.next(),
result, index);
+ throw e;
}
return result;
}
private int doCreateProcessorExchangePairs(
- Exchange exchange, Object recipient, List<ProcessorExchangePair>
result, int index)
+ Exchange exchange, Object recipient, List<ProcessorExchangePair>
result, int index,
+ Map<String, Object> txData)
throws NoTypeConversionAvailableException {
boolean prototype = cacheSize < 0;
@@ -330,7 +326,7 @@ public class RecipientListProcessor extends
MulticastProcessor {
}
// then create the exchange pair
- result.add(createProcessorExchangePair(index++, endpoint, producer,
exchange, pattern, prototype));
+ result.add(createProcessorExchangePair(index++, endpoint, producer,
exchange, pattern, prototype, txData));
return index;
}
@@ -340,6 +336,13 @@ public class RecipientListProcessor extends
MulticastProcessor {
protected ProcessorExchangePair createProcessorExchangePair(
int index, Endpoint endpoint, Producer producer,
Exchange exchange, ExchangePattern pattern, boolean
prototypeEndpoint) {
+ return createProcessorExchangePair(index, endpoint, producer,
exchange, pattern, prototypeEndpoint,
+ exchange.isTransacted() ? new ConcurrentHashMap<>() : null);
+ }
+
+ private ProcessorExchangePair createProcessorExchangePair(
+ int index, Endpoint endpoint, Producer producer,
+ Exchange exchange, ExchangePattern pattern, boolean
prototypeEndpoint, Map<String, Object> txData) {
// copy exchange, and do not share the unit of work
Exchange copy =
processorExchangeFactory.createCorrelatedCopy(exchange, false);
copy.getExchangeExtension().setTransacted(exchange.isTransacted());
@@ -350,10 +353,7 @@ public class RecipientListProcessor extends
MulticastProcessor {
// If we are in a transaction, set TRANSACTION_CONTEXT_DATA property
for new exchanges to share txData
// during the transaction.
- if (exchange.isTransacted() &&
copy.getProperty(Exchange.TRANSACTION_CONTEXT_DATA) == null) {
- if (txData == null) {
- txData = new ConcurrentHashMap<>();
- }
+ if (txData != null &&
copy.getProperty(Exchange.TRANSACTION_CONTEXT_DATA) == null) {
copy.setProperty(Exchange.TRANSACTION_CONTEXT_DATA, txData);
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/Splitter.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/Splitter.java
index 050872386d75..c00f2c13045f 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/Splitter.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/Splitter.java
@@ -292,8 +292,7 @@ public class Splitter extends MulticastProcessor {
// tracks individual (raw) item count, independent of grouping
private final AtomicInteger rawItemCount = new AtomicInteger();
// tracks whether the primary (processing) iterator has been created;
- // subsequent iterators (e.g. the drain in MulticastProcessor.doDone)
- // must not update the watermark count (CAMEL-24139)
+ // any subsequent iterator must not update the watermark count
(CAMEL-24139)
private boolean primaryIteratorCreated;
private SplitterIterable(Exchange exchange, Object value) {
@@ -363,7 +362,7 @@ public class Splitter extends MulticastProcessor {
@Override
public Iterator<ProcessorExchangePair> iterator() {
// only the first (primary) iterator tracks watermark count;
- // subsequent iterators (drain in doDone) must not inflate it
(CAMEL-24139)
+ // subsequent iterators must not inflate it (CAMEL-24139)
boolean isPrimary = !primaryIteratorCreated;
primaryIteratorCreated = true;
@@ -440,7 +439,7 @@ public class Splitter extends MulticastProcessor {
}
}
// eagerly update watermark count for items actually
routed (primary iterator only)
- // so the drain loop in MulticastProcessor.doDone
cannot inflate it (CAMEL-24139)
+ // so a subsequent iterator cannot inflate it
(CAMEL-24139)
if (isPrimary && resumeStrategy != null &&
watermarkKey != null && watermarkExpression == null) {
original.setProperty(SPLIT_WATERMARK_COUNT,
rawItemCount.get());
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListInvalidEndpointReleaseTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListInvalidEndpointReleaseTest.java
new file mode 100644
index 000000000000..688bc84df331
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListInvalidEndpointReleaseTest.java
@@ -0,0 +1,117 @@
+/*
+ * 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;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.DefaultProducer;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * When a recipient cannot be resolved, the producers already acquired for the
recipients before it must be released.
+ */
+public class RecipientListInvalidEndpointReleaseTest extends
ContextTestSupport {
+
+ private final Map<String, AtomicInteger> started = new
ConcurrentHashMap<>();
+ private final Map<String, AtomicInteger> stopped = new
ConcurrentHashMap<>();
+
+ @Test
+ public void testProducersReleasedWhenLaterRecipientIsInvalid() {
+ assertThrows(Exception.class, () -> template.sendBody("direct:start",
"Hello"));
+
+ assertEquals(1, count(started, "a"));
+ assertEquals(1, count(stopped, "a"), "the producer of the prototype
recipient should be released and stopped");
+ }
+
+ @Test
+ public void testCachedProducerReleasedWhenLaterRecipientIsInvalid() {
+ assertThrows(Exception.class, () -> template.sendBody("direct:cached",
"Hello"));
+ assertThrows(Exception.class, () -> template.sendBody("direct:cached",
"World"));
+
+ // the producer is not a singleton, so it is pooled: when released
after the first message,
+ // the second message acquires it from the pool again, instead of
creating a new producer
+ assertEquals(1, count(started, "c"), "the pooled producer should be
released and reused");
+ }
+
+ private static int count(Map<String, AtomicInteger> counters, String key) {
+ AtomicInteger counter = counters.get(key);
+ return counter != null ? counter.get() : 0;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ context.addComponent("track", new DefaultComponent() {
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ return new DefaultEndpoint(uri, this) {
+ @Override
+ public Producer createProducer() {
+ return new DefaultProducer(this) {
+ @Override
+ public void process(Exchange exchange) {
+ // noop
+ }
+
+ @Override
+ protected void doStart() {
+ started.computeIfAbsent(remaining, k -> new
AtomicInteger()).incrementAndGet();
+ }
+
+ @Override
+ protected void doStop() {
+ stopped.computeIfAbsent(remaining, k -> new
AtomicInteger()).incrementAndGet();
+ }
+ };
+ }
+
+ @Override
+ public boolean isSingletonProducer() {
+ return false;
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ throw new UnsupportedOperationException();
+ }
+ };
+ }
+ });
+
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").recipientList(constant("track:a,unknownxyz:b")).cacheSize(-1);
+
+
from("direct:cached").recipientList(constant("track:c,unknownxyz:d"));
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListTransactedContextDataTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListTransactedContextDataTest.java
new file mode 100644
index 000000000000..b7124554d739
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/RecipientListTransactedContextDataTest.java
@@ -0,0 +1,81 @@
+/*
+ * 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;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+/**
+ * Each transaction must get its own transaction context data when a
transacted exchange goes through a recipient list.
+ */
+public class RecipientListTransactedContextDataTest extends ContextTestSupport
{
+
+ private final List<Map<?, ?>> data = new ArrayList<>();
+
+ @Test
+ public void testTransactionContextDataIsNotShared() {
+ template.sendBody("direct:start", "A");
+ template.sendBody("direct:start", "B");
+
+ assertEquals(2, data.size());
+ assertNotNull(data.get(0));
+ assertNotNull(data.get(1));
+ assertNotSame(data.get(0), data.get(1), "two transactions must not
share the transaction context data");
+ }
+
+ @Test
+ public void testTransactionContextDataIsSharedWithinTransaction() {
+ template.sendBody("direct:two", "A");
+ template.sendBody("direct:two", "B");
+
+ assertEquals(4, data.size());
+ assertNotNull(data.get(0));
+ assertSame(data.get(0), data.get(1), "the recipients of one
transaction must share the transaction context data");
+ assertSame(data.get(2), data.get(3), "the recipients of one
transaction must share the transaction context data");
+ assertNotSame(data.get(0), data.get(2), "two transactions must not
share the transaction context data");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .process(e ->
e.getExchangeExtension().setTransacted(true))
+ .recipientList(constant("direct:a"));
+
+ from("direct:two")
+ .process(e ->
e.getExchangeExtension().setTransacted(true))
+ .recipientList(constant("direct:a,direct:a"));
+
+ from("direct:a")
+ .process(e ->
data.add(e.getProperty(Exchange.TRANSACTION_CONTEXT_DATA, Map.class)));
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/SplitterStreamingStopOnExceptionReadAheadTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/SplitterStreamingStopOnExceptionReadAheadTest.java
new file mode 100644
index 000000000000..349f1d7ffee5
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/SplitterStreamingStopOnExceptionReadAheadTest.java
@@ -0,0 +1,78 @@
+/*
+ * 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;
+
+import java.util.Iterator;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * A streaming split that stops on an exception must not read, prepare or
create exchanges for the parts after the
+ * failure.
+ */
+public class SplitterStreamingStopOnExceptionReadAheadTest extends
ContextTestSupport {
+
+ private final AtomicInteger read = new AtomicInteger();
+ private final AtomicInteger prepared = new AtomicInteger();
+
+ @Test
+ public void testStopOnExceptionDoesNotReadRemainingParts() {
+ Iterator<Integer> parts = new Iterator<>() {
+ private int next;
+
+ @Override
+ public boolean hasNext() {
+ return next < 100;
+ }
+
+ @Override
+ public Integer next() {
+ read.incrementAndGet();
+ return next++;
+ }
+ };
+
+ assertThrows(Exception.class, () -> template.sendBody("direct:start",
parts));
+
+ // part 0 is routed, part 1 fails and stops the split
+ assertEquals(2, prepared.get(), "onPrepare should only be called for
the parts that are routed");
+ assertEquals(2, read.get(), "the parts after the failure should not be
read");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+
.split(body()).streaming().stopOnException().onPrepare(e ->
prepared.incrementAndGet())
+ .process(e -> {
+ if (e.getMessage().getBody(Integer.class) == 1) {
+ throw new IllegalArgumentException("Forced");
+ }
+ })
+ .end();
+ }
+ };
+ }
+}