You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Camel开启parallelAggregate后聚合方法未并行执行问题求助

看起来你已经尝试开启了parallelAggregate来让Split的聚合操作并行,但实际还是串行执行——这个问题我之前也碰到过,根源在于Camel内部的AsyncCompletionService用了ReentrantLock来保护结果队列,哪怕开了并行聚合,所有聚合操作还是会被这个锁给串行化了。

从你的线程dump里也能看出来:一个Split线程在等待获取这个锁(parking to wait for <0x000000076e105c70>),而另一个线程正卡在你的aggregate方法的sleep里,还持有这个锁——这就导致所有聚合请求只能排队等着,一个执行完下一个才上。

接下来给你几个解决方向:

1. 先确保你的聚合策略是线程安全的

你的当前聚合逻辑只是返回第一个完成的Exchange,但如果之后要做结果合并(比如收集所有子任务的输出),必须用线程安全的容器来存数据,比如CopyOnWriteArrayList或者ConcurrentHashMap,不然并行聚合会出现数据竞争的问题。

2. 用AsyncAggregationStrategy实现真正的异步聚合

Camel的AsyncAggregationStrategy允许你把聚合逻辑放到异步线程里执行,这样即使AsyncCompletionService的锁让提交操作串行,实际的聚合逻辑可以在多个线程里并行跑。给你改个示例:

@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 {
                // 自定义异步聚合策略
                AsyncAggregationStrategy asyncAggStrategy = new AsyncAggregationStrategy() {
                    // 同步聚合逻辑(线程安全的实现)
                    @Override
                    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
                        try {
                            Thread.sleep(1000);
                        } catch (InterruptedException e) {
                            LOGGER.error("InterruptedException: ", e);
                        }
                        // 这里根据你的业务调整,比如收集所有结果
                        if (oldExchange == null) {
                            List<String> results = new ArrayList<>();
                            results.add(newExchange.getIn().getBody(String.class));
                            newExchange.getIn().setBody(results);
                            return newExchange;
                        } else {
                            List<String> results = oldExchange.getIn().getBody(List.class);
                            results.add(newExchange.getIn().getBody(String.class));
                            return oldExchange;
                        }
                    }

                    // 异步执行聚合
                    @Override
                    public void aggregate(Exchange oldExchange, Exchange newExchange, AsyncCallback callback) {
                        // 用自定义线程池执行聚合逻辑
                        Executors.newFixedThreadPool(4).submit(() -> {
                            try {
                                aggregate(oldExchange, newExchange);
                            } finally {
                                // 通知Camel聚合完成
                                callback.done(false);
                            }
                        });
                    }
                };

                from("direct:start")
                    .split(body(), asyncAggStrategy)
                    .streaming()
                    .parallelProcessing()
                    .parallelAggregate()
                    .executorService(Executors.newFixedThreadPool(4)) // 指定并行处理的线程池
                    .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);
    }
}

3. 调整线程池配置

别忘了给Split指定足够大小的线程池,不然线程数不够的话,并行效果也出不来。上面的示例里用了executorService(Executors.newFixedThreadPool(4)),你可以根据自己的任务量调整大小。

另外要注意:parallelAggregate的核心是允许子任务完成后立即聚合,而不是等所有子任务都结束再统一聚合。如果你的聚合逻辑依赖之前的聚合结果(比如累积计算),并行聚合可能会导致结果顺序不可控,这个得根据你的业务场景来权衡。

内容的提问来源于stack exchange,提问作者Jey350

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.27 16:52:35