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

Camel拆分后聚合场景下的异常传播问题:如何终止后续处理

Apache Camel Split-Aggregate 异常传播解决方案

问题根源

当前代码中,聚合完成后抛出的异常无法被聚合策略捕获,且最后一步的process仍会执行,原因在于:

  • 并行处理(parallelProcessing)模式下,Split组件默认不共享工作单元,子Exchange的异常无法传递到父Exchange;
  • throwException抛出的异常属于聚合完成后的处理器环节,不会被回传给聚合策略的newExchange参数;
  • stopOnException在并行模式下需要额外配置才能生效。

修复方案

关键改动点

  • 给Split添加shareUnitOfWork(),确保并行处理时异常能传递到父Exchange;
  • 在聚合后的异常抛出环节,确保Exchange被标记为异常状态;
  • 调整聚合策略,当捕获到异常Exchange时,直接返回异常Exchange,终止后续聚合与流程。

修改后的完整代码

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

    @Override
    protected RouteBuilder createRouteBuilder() throws Exception {
        return new RouteBuilder() {
            @Override
            public void configure() throws Exception {
                // 全局异常处理:确保异常被正确传递,不提前拦截
                onException(Exception.class)
                        .log("捕获异常: ${exception.message}")
                        .handled(false);

                from("direct:start")    
                    .split(body())
                    .streaming()
                    .stopOnException()
                    .parallelProcessing()
                    .shareUnitOfWork() // 核心配置:并行模式下共享工作单元,传递异常到父Exchange
                        .aggregate(new AggregationStrategy() {
                            @Override
                            public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
                                // 检查子Exchange是否携带异常
                                if (newExchange.getException() != null) {
                                    LOGGER.info("聚合策略捕获到异常");
                                    // 直接返回异常Exchange,终止聚合流程
                                    return newExchange;
                                }
                                // 正常聚合逻辑
                                return oldExchange == null ? newExchange : oldExchange;
                            }
                        }).constant(true)
                        .completionSize(1)
                        .completionTimeout(500)
                            .log(LoggingLevel.INFO, LOGGER, "Aggreg ${body}")
                            .throwException(Exception.class, "propagate plz")
                        .end()
                    .end()
                    // 仅当无异常时才会执行此处理器
                    .process(e -> {
                        LOGGER.info("I don't want to be seen, because of {}", e.getException());
                    });
            }
        };
    }
    
    @Test
    public void test1() throws InterruptedException {
        template.sendBody("direct:start", Arrays.asList("A", "B", "C", "D"));
        Thread.sleep(5000);
    }
}

额外说明

  • shareUnitOfWork()是并行模式下异常传播的核心,它让所有并行子任务共享同一工作单元,确保异常能向上传递到父Exchange;
  • onException设置handled(false)是为了让异常继续传播,不被提前拦截,确保Split的stopOnException能触发终止后续子任务;
  • 聚合策略中捕获到异常Exchange后直接返回,会终止整个聚合流程,父Exchange携带异常后,最后一步的process就不会执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 11:36:06