如何在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
相关产品推荐
相关产品推荐

