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

