Camel parallelAggregate

Viewed 84

I'll need the aggregate method to be run in parallel in my split so I activated parallelAggregate option but it doesn't work as I expected it.

@RunWith(SpringJUnit4ClassRunner.class)
public class SplitAggregationTest extends CamelTestSupport {
private final Logger LOGGER = LoggerFactory.getLogger(SplitAggregationTest.class);

@Override
protected RouteBuilder createRouteBuilder() throws Exception {
    return new RouteBuilder() {
        @Override
        public void configure() throws Exception {
            from("direct:start")
                    .split(body(), new AggregationStrategy() {
                        @Override
                        public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
                            try {
                                Thread.sleep(1000);
                            } catch (InterruptedException e) {
                                LOGGER.error("InterruptedException: ", e);
                            }
                            return oldExchange == null ? newExchange : oldExchange;
                        }
                    }).streaming().parallelProcessing().parallelAggregate()
                        .log(LoggingLevel.INFO, LOGGER, "Aggreg ${body}")
                    .end();
        }
    };
}

@Test
public void test1() throws InterruptedException {
    Thread.sleep(5000);

    template.sendBody("direct:start", Arrays.asList("A", "B", "C", "D"));

    Thread.sleep(5000);
}

Aggregate methods (in purple) are not run in parallel, they are waiting each other to finish before:

Threads view

Edit - The according thread dump :

"Camel (camel-2) thread #17 - Split" #36 daemon prio=5 os_prio=0   tid=0x000002317da51800 nid=0x149bc waiting on condition  [0x00000062a55fe000]
java.lang.Thread.State: WAITING (parking)
at sun.misc.Unsafe.park(Native Method)
- parking to wait for  <0x000000076e105c70> (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:836)
at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireQueued(AbstractQueuedSynchronizer.java:870)
at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:1199)
at java.util.concurrent.locks.ReentrantLock$NonfairSync.lock(ReentrantLock.java:209)
at java.util.concurrent.locks.ReentrantLock.lock(ReentrantLock.java:285)
at org.apache.camel.util.concurrent.AsyncCompletionService.complete(AsyncCompletionService.java:141)
at org.apache.camel.util.concurrent.AsyncCompletionService.access$200(AsyncCompletionService.java:30)
at org.apache.camel.util.concurrent.AsyncCompletionService$Task.accept(AsyncCompletionService.java:168)
at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask.lambda$null$0(MulticastProcessor.java:579)
at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask$$Lambda$519/1049021104.done(Unknown Source)
at org.apache.camel.AsyncCallback.run(AsyncCallback.java:44)
at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187)
at org.apache.camel.impl.engine.DefaultReactiveExecutor.schedule(DefaultReactiveExecutor.java:59)
at org.apache.camel.processor.MulticastProcessor.lambda$schedule$1(MulticastProcessor.java:348)
at org.apache.camel.processor.MulticastProcessor$$Lambda$517/1626228981.run(Unknown Source)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

Locked ownable synchronizers:
- <0x000000076f0ca838> (a java.util.concurrent.ThreadPoolExecutor$Worker)

"Camel (camel-2) thread #16 - Split" #35 daemon prio=5 os_prio=0   tid=0x000002317da54000 nid=0x6a8 waiting on condition [0x00000062a54fe000]
java.lang.Thread.State: TIMED_WAITING (sleeping)
at java.lang.Thread.sleep(Native Method)
at com.mybatch.camel.processor.SplitAggregationTest$1$1.aggregate(SplitAggregationTest.java:34)
at org.apache.camel.AggregationStrategy.aggregate(AggregationStrategy.java:86)
at org.apache.camel.processor.MulticastProcessor.doAggregateInternal(MulticastProcessor.java:894)
at org.apache.camel.processor.MulticastProcessor.doAggregate(MulticastProcessor.java:858)
at org.apache.camel.processor.MulticastProcessor$MulticastTask.aggregate(MulticastProcessor.java:451)
at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask.lambda$null$0(MulticastProcessor.java:582)
at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask$$Lambda$519/1049021104.done(Unknown Source)
at org.apache.camel.AsyncCallback.run(AsyncCallback.java:44)
at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187)
at org.apache.camel.impl.engine.DefaultReactiveExecutor.schedule(DefaultReactiveExecutor.java:59)
at org.apache.camel.processor.MulticastProcessor.lambda$schedule$1(MulticastProcessor.java:348)
at org.apache.camel.processor.MulticastProcessor$$Lambda$517/1626228981.run(Unknown Source)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

Locked ownable synchronizers:
- <0x000000076e105c70> (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
- <0x000000076ef115b0> (a java.util.concurrent.ThreadPoolExecutor$Worker)
2 Answers

Let's define these terms first:

parallelProcessing: If true, Each split item will be run via threads from a thread pool, either default, or one you provide.

parallelAggregate: If true, Camel will not synchronize calls to your AggregationStrategy (make sure it's thread safe btw).

Based on the names of the threads you're getting, Camel is running the split items on different threads. Are you certain that your aggregate method is not being called concurrently by multiple threads? Perhaps what you're seeing is that Camel waits for all split items to complete before moving past the split.

After some investigations (src here) it appears the reentrantlock in MulticastProcessor#agregate() prevents the split threads to process the AggregationStrategy at the same time from camel V3.X (in V2.X the same test aggregates as expected in // )

Don't you think this reentrantlock in agregate methode is not in the sense of parallelAggregate config ?

Here is the log output of the test (agregate processing during 1 sec, one after other)

2022-08-12 17:09:14,997|INFO |main|o.a.c.i.e.AbstractCamelContext:3116||    Started route1 (direct://start)
2022-08-12 17:09:14,997|INFO |main|o.a.c.i.e.AbstractCamelContext:3134||Apache Camel 3.14.0 (camel-1) started in 169ms (build:48ms init:114ms start:7ms)
2022-08-12 17:09:15,013|INFO |pool-1-thread-3|SplitAggregationTest:166||Split part C
2022-08-12 17:09:15,013|INFO |pool-1-thread-2|SplitAggregationTest:166||Split part A
2022-08-12 17:09:15,013|INFO |pool-1-thread-1|SplitAggregationTest:166||Split part D
2022-08-12 17:09:15,013|INFO |pool-1-thread-4|SplitAggregationTest:166||Split part B
2022-08-12 17:09:15,016|INFO |pool-1-thread-4|SplitAggregationTest:29||> Aggregate NULL to new: B:FCFE9D041FBDD72-0000000000000001
2022-08-12 17:09:16,016|INFO |pool-1-thread-1|SplitAggregationTest:29||> Aggregate B:FCFE9D041FBDD72-0000000000000001 to new: D:FCFE9D041FBDD72-0000000000000002
2022-08-12 17:09:17,016|INFO |pool-1-thread-3|SplitAggregationTest:29||> Aggregate B:FCFE9D041FBDD72-0000000000000001 to new: C:FCFE9D041FBDD72-0000000000000003
2022-08-12 17:09:18,016|INFO |pool-1-thread-2|SplitAggregationTest:29||> Aggregate B:FCFE9D041FBDD72-0000000000000001 to new: A:FCFE9D041FBDD72-0000000000000004
2022-08-12 17:09:19,016|INFO |pool-1-thread-2|SplitAggregationTest:166||End aggregation  B

Here the threaddump when "pool-1-thread-1" is agregating. the others 3 threads are waiting parking to wait for ReentrantLock$NonfairSync

"pool-1-thread-4" #15 prio=5 os_prio=0 tid=0x00000000210a1800 nid=0x5988 waiting on condition [0x000000002275e000]
   java.lang.Thread.State: WAITING (parking)
        at sun.misc.Unsafe.park(Native Method)
        - parking to wait for  <0x000000076bd49968> (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
        at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:836)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireQueued(AbstractQueuedSynchronizer.java:870)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:1199)
        at java.util.concurrent.locks.ReentrantLock$NonfairSync.lock(ReentrantLock.java:209)
        at java.util.concurrent.locks.ReentrantLock.lock(ReentrantLock.java:285)
        at org.apache.camel.util.concurrent.AsyncCompletionService.complete(AsyncCompletionService.java:141)
        at org.apache.camel.util.concurrent.AsyncCompletionService.access$200(AsyncCompletionService.java:30)
        at org.apache.camel.util.concurrent.AsyncCompletionService$Task.accept(AsyncCompletionService.java:168)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask.lambda$null$0(MulticastProcessor.java:579)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask$$Lambda$73/792205377.done(Unknown Source)
        at org.apache.camel.AsyncCallback.run(AsyncCallback.java:44)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor.schedule(DefaultReactiveExecutor.java:59)
        at org.apache.camel.processor.MulticastProcessor.lambda$schedule$1(MulticastProcessor.java:348)
        at org.apache.camel.processor.MulticastProcessor$$Lambda$70/900268191.run(Unknown Source)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

   Locked ownable synchronizers:
        - <0x000000076f701340> (a java.util.concurrent.ThreadPoolExecutor$Worker)

"pool-1-thread-3" #14 prio=5 os_prio=0 tid=0x00000000210a0000 nid=0x3c8c waiting on condition [0x000000002265e000]
   java.lang.Thread.State: WAITING (parking)
        at sun.misc.Unsafe.park(Native Method)
        - parking to wait for  <0x000000076bd49968> (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
        at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:836)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireQueued(AbstractQueuedSynchronizer.java:870)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:1199)
        at java.util.concurrent.locks.ReentrantLock$NonfairSync.lock(ReentrantLock.java:209)
        at java.util.concurrent.locks.ReentrantLock.lock(ReentrantLock.java:285)
        at org.apache.camel.util.concurrent.AsyncCompletionService.complete(AsyncCompletionService.java:141)
        at org.apache.camel.util.concurrent.AsyncCompletionService.access$200(AsyncCompletionService.java:30)
        at org.apache.camel.util.concurrent.AsyncCompletionService$Task.accept(AsyncCompletionService.java:168)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask.lambda$null$0(MulticastProcessor.java:579)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask$$Lambda$73/792205377.done(Unknown Source)
        at org.apache.camel.AsyncCallback.run(AsyncCallback.java:44)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor.schedule(DefaultReactiveExecutor.java:59)
        at org.apache.camel.processor.MulticastProcessor.lambda$schedule$1(MulticastProcessor.java:348)
        at org.apache.camel.processor.MulticastProcessor$$Lambda$70/900268191.run(Unknown Source)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

   Locked ownable synchronizers:
        - <0x000000076f7011d0> (a java.util.concurrent.ThreadPoolExecutor$Worker)

"pool-1-thread-2" #13 prio=5 os_prio=0 tid=0x00000000210a0800 nid=0x32f8 waiting on condition [0x000000002255f000]
   java.lang.Thread.State: WAITING (parking)
        at sun.misc.Unsafe.park(Native Method)
        - parking to wait for  <0x000000076f7010f8> (a java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject)
        at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2039)
        at java.util.concurrent.LinkedBlockingQueue.take(LinkedBlockingQueue.java:442)
        at java.util.concurrent.ThreadPoolExecutor.getTask(ThreadPoolExecutor.java:1074)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1134)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

   Locked ownable synchronizers:
        - None

"pool-1-thread-1" #12 prio=5 os_prio=0 tid=0x000000002109f000 nid=0x4770 runnable [0x000000002245e000]
   java.lang.Thread.State: RUNNABLE
        at SplitAggregationTest$1.lambda$0(SplitAggregationTest.java:35)
        at SplitAggregationTest$1$$Lambda$30/2021707251.aggregate(Unknown Source)
        at org.apache.camel.AggregationStrategy.aggregate(AggregationStrategy.java:86)
        at org.apache.camel.processor.MulticastProcessor.doAggregateInternal(MulticastProcessor.java:894)
        at org.apache.camel.processor.MulticastProcessor.doAggregate(MulticastProcessor.java:858)
        at org.apache.camel.processor.MulticastProcessor$MulticastTask.aggregate(MulticastProcessor.java:451)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask.lambda$null$0(MulticastProcessor.java:582)
        at org.apache.camel.processor.MulticastProcessor$MulticastReactiveTask$$Lambda$73/792205377.done(Unknown Source)
        at org.apache.camel.AsyncCallback.run(AsyncCallback.java:44)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187)
        at org.apache.camel.impl.engine.DefaultReactiveExecutor.schedule(DefaultReactiveExecutor.java:59)
        at org.apache.camel.processor.MulticastProcessor.lambda$schedule$1(MulticastProcessor.java:348)
        at org.apache.camel.processor.MulticastProcessor$$Lambda$70/900268191.run(Unknown Source)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

   Locked ownable synchronizers:
        - <0x000000076bd49968> (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
        - <0x000000076f701390> (a java.util.concurrent.ThreadPoolExecutor$Worker)

Thanx for any help

Related