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

如何在Apache Camel中为特定SEDA消费者条件化使用onCompletion()

问题

我有多个消费者连接到SEDA队列,每个消费者通过choice-when自行决定是否处理事件。每个消费者需要在后续所有路由的整个流程完成或失败时做出响应。

原代码(存在于父类并被多次继承)如下:

//this exists in a super class and is inherited multiple times
from("seda:queue-with-multiple-consumers?multipleConsumers=true")
    .routeId("id-of-some-consumer-1") //will be set per consumer
    .choice().when(consumerCheckIfEligible())
        .onCompletion().onCompleteOnly()
            .process(exchange -> {
                //should be triggered if direct:delegate-to-actual-processing-route succeeds
                log.info("onCompletion().onCompleteOnly()");
            })
        .end()
        
        .onCompletion().onFailureOnly()
            .process(exchange -> {
                //should be triggered if direct:delegate-to-actual-processing-route fails
                log.info("onCompletion().onFailureOnly()");
            })
        .end()

        .onException(Exception.class)
            .process(e -> {
                //should be triggered if direct:delegate-to-actual-processing-route fails
                log.info("onException(Exception.class)");
            })
        .end()

        .to("direct:delegate-to-actual-processing-route")

    .end()    

但Apache Camel不允许onCompletion()钩子放在非顶层位置,报错如下:

java.lang.IllegalArgumentException: The output must be added as top-level on the route. Try moving onCompletion[[]] to the top of route.
    at org.apache.camel.spring.boot.CamelSpringBootApplicationListener.onApplicationEvent(CamelSpringBootApplicationListener.java:212)
    ...(省略部分栈信息)

需绕过该报错,实现条件化使用onCompletion()/onException()钩子。

解决方案

方法1:将onCompletion移至路由顶层,通过Exchange属性做条件判断

把onCompletion和onException放在路由最外层,先在choice分支中给符合处理条件的Exchange标记专属属性,后续钩子逻辑通过过滤该属性来实现条件触发。

修改后代码示例:

from("seda:queue-with-multiple-consumers?multipleConsumers=true")
    .routeId("id-of-some-consumer-1")
    // 标记符合条件的Exchange
    .choice()
        .when(consumerCheckIfEligible())
            .setProperty("eligibleForProcessing", constant(true))
            .to("direct:delegate-to-actual-processing-route")
        .otherwise()
            // 不符合条件直接终止流程
            .stop()
    .end()
    // 顶层成功完成钩子:仅处理标记过的Exchange
    .onCompletion().onCompleteOnly()
        .filter(exchangeProperty("eligibleForProcessing").isEqualTo(true))
        .process(exchange -> {
            log.info("onCompletion().onCompleteOnly()");
        })
    .end()
    // 顶层失败完成钩子:仅处理标记过的Exchange
    .onCompletion().onFailureOnly()
        .filter(exchangeProperty("eligibleForProcessing").isEqualTo(true))
        .process(exchange -> {
            log.info("onCompletion().onFailureOnly()");
        })
    .end()
    // 顶层异常处理:仅处理标记过的Exchange
    .onException(Exception.class)
        .filter(exchangeProperty("eligibleForProcessing").isEqualTo(true))
        .process(e -> {
            log.info("onException(Exception.class)");
        })
    .end();

方法2:将条件处理逻辑封装为子路由

把符合条件后的完整处理流程(包括钩子)抽离成独立的direct子路由,原路由的choice分支直接转发到该子路由。子路由属于独立顶层路由,可正常使用onCompletion和onException。

父类主路由:

from("seda:queue-with-multiple-consumers?multipleConsumers=true")
    .routeId("id-of-some-consumer-1")
    .choice()
        .when(consumerCheckIfEligible())
            .to("direct:consumer-specific-processing")
    .end();

子路由(可在子类中定义或继承复用):

from("direct:consumer-specific-processing")
    .onCompletion().onCompleteOnly()
        .process(exchange -> {
            log.info("onCompletion().onCompleteOnly()");
        })
    .end()
    .onCompletion().onFailureOnly()
        .process(exchange -> {
            log.info("onCompletion().onFailureOnly()");
        })
    .end()
    .onException(Exception.class)
        .process(e -> {
            log.info("onException(Exception.class)");
        })
    .end()
    .to("direct:delegate-to-actual-processing-route");

方法3:自定义处理器或路由策略(进阶场景)

如果需要更细粒度的控制,可自定义处理器在delegate-to-actual-processing-route执行完成后,手动判断结果并触发响应逻辑;或者使用RoutePolicy监听路由生命周期,在Exchange完成/失败时根据条件执行对应逻辑。这种方式复杂度较高,适合特殊业务场景。

内容的提问来源于stack exchange,提问作者Marian Klühspies

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 06:15:33