架构

OnCompletion

师成师成· 更新于 2026-09-28· 阅读 18 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

OnCompletion

Camel 中有一个涵盖 Exchange 的 Unit of Work(工作单元)概念。工作单元除其他功能外,还支持在 Exchange 完成时调用的同步回调。回调 API 定义在 org.apache.camel.spi.Synchronization 中,同时还提供了扩展的同步接口 org.apache.camel.spi.SynchronizationRouteAware,用于路由事件的回调。

UnitOfWork API

可以通过 org.apache.camel.Exchange 的 getUnitOfWork() 方法获取 org.apache.camel.spi.UnitOfWork。

OnCompletion DSL

OnCompletion EIP 支持以下特性:

  • scope:context 或 route(优先级高于全局设置)
  • trigger:always,或者仅在成功或失败时触发
  • onWhen:谓词,仅在匹配该谓词时才触发
  • mode:定义是在路由消费者将响应写回调用方之前还是之后运行(如果为 InOut 模式)(默认为 AfterConsumer)
  • parallelProcessing:是异步还是同步运行(是否使用线程池)(默认为 false)

onCompletion 支持以同步或异步模式(使用线程池)运行完成任务,也支持在路由消费者完成之前或之后运行。这样设计是为了提供更多灵活性。例如,可以指定在路由消费者完成之前同步运行,这样就能在消费者向调用方写回任何响应之前修改 Exchange。你可以利用这一点添加自定义头信息、将响应消息发送到日志进行记录等。

带路由作用域的 onCompletion

OnCompletion EIP 允许你在原始 Exchange 完成时添加自定义路由/处理器。Camel 会复制一份 Exchange,并在单独的线程中对其进行路由,这与 Wire Tap 类似。这样,原始线程可以在 onCompletion 路由并发运行的同时继续执行。我们选择这种模型,是因为不希望 onCompletion 路由对原始路由产生干扰。

多个 onCompletion

你可以在 context 和 route 两个层级上定义多个 onCompletion。

当你定义了路由作用域的 onCompletion 后,该路由上所有 context 作用域的 onCompletion 都会被禁用。

  • Java
  • XML
  • YAML
from("direct:start")
    .onCompletion()
        // this route is only invoked when the original route is complete as a kind
        // of completion callback
        .to("log:sync")
        .to("mock:sync")
    // must use end to denote the end of the onCompletion route
    .end()
    // here the original route contiunes
    .process(new MyProcessor())
    .to("mock:result");
<route>
    <from uri="direct:start"/>
    <!-- this onCompletion block will only be executed when the exchange is done being routed -->
    <!-- this callback is always triggered even if the exchange failed -->
    <onCompletion>
        <!-- so this is a kinda like an after completion callback -->
        <to uri="log:sync"/>
        <to uri="mock:sync"/>
    </onCompletion>
    <process ref="myProcessor"/>
    <to uri="mock:result"/>
</route>

在 YAML 中,路由作用域的 onCompletion 在路由所引用的 routeConfiguration 中声明,因此它只对该路由生效。

- routeConfiguration:
    id: myConfig
    onCompletion:
      - onCompletion:
          steps:
            - to:
                uri: log:sync
            - to:
                uri: mock:sync
- route:
    routeConfigurationId: myConfig
    from:
      uri: direct:start
      steps:
        - process:
            ref: myProcessor
        - to:
            uri: mock:result

默认情况下,OnCompletion EIP 会在 Exchange 完成时始终被触发。如果你只希望在交换成功完成时触发,或只在失败时触发,则可以通过 onCompleteOnly 或 onFailureOnly 来指定,如下所示:

  • Java
  • XML
  • YAML
from("direct:start")
    // here we qualify onCompletion to only invoke when the exchange failed (exception or FAULT body)
    .onCompletion().onFailureOnly()
        .to("log:sync")
        .to("mock:syncFail")
    // must use end to denote the end of the onCompletion route
    .end()
    .onCompletion().onCompleteOnly()
        .to("log:sync")
        .to("mock:syncOK")
    .end()
    // here the original route continues
    .process(new MyProcessor())
    .to("mock:result");
<route>
    <from uri="direct:start"/>
    <!-- this onCompletion block will only be executed when the exchange is done being routed -->
    <!-- this callback is only triggered when the exchange failed, as we have onFailureOnly=true -->
    <onCompletion onFailureOnly="true">
        <to uri="log:sync"/>
        <to uri="mock:sync"/>
    </onCompletion>
    <process ref="myProcessor"/>
    <to uri="mock:result"/>
</route>
- routeConfiguration:
    id: myConfig
    onCompletion:
      - onCompletion:
          onFailureOnly: "true"
          steps:
            - to:
                uri: log:sync
            - to:
                uri: mock:sync
- route:
    routeConfigurationId: myConfig
    from:
      uri: direct:start
      steps:
        - process:
            ref: myProcessor
        - to:
            uri: mock:result

你可以通过 Exchange.ON_COMPLETION 属性来判断一个 Exchange 是否为 OnCompletion Exchange,Camel 会为该属性设置布尔值 true。

全局级别的 onCompletion

其工作方式与路由级别相同,只是它们是全局定义的。示例如下:

  • Java
  • XML
  • YAML
// define a global on completion that is invoked when the exchange is done being routed
onCompletion().to("log:global").to("mock:sync");

from("direct:start")
    .process(new MyProcessor())
    .to("mock:result");
<!-- this is a global onCompletion route that is invoked when any exchange is done being routed
     as a kind of after callback -->
<onCompletion>
    <to uri="log:global"/>
    <to uri="mock:sync"/>
</onCompletion>

<route>
    <from uri="direct:start"/>
    <process ref="myProcessor"/>
    <to uri="mock:result"/>
</route>
- onCompletion:
    steps:
      - to:
          uri: log:global
      - to:
          uri: mock:sync
- route:
    from:
      uri: direct:start
      steps:
        - process:
            ref: myProcessor
        - to:
            uri: mock:result
如果在路由中定义了 onCompletion,它会覆盖所有全局作用域的配置,因此只会使用路由作用域的配置。全局作用域的配置将不再生效。

在 onCompletion 中使用 onWhen 谓词

与 Camel 中的其他 DSL 一样,你可以为 onCompletion 附加一个谓词,这样它只会在谓词匹配、即满足特定条件时才触发。例如,要让它仅在消息体包含 Hello 一词时触发,可以这样写:

  • Java
  • XML
  • YAML
from("direct:start")
    .onCompletion().onWhen(body().contains("Hello"))
        // this route is only invoked when the original route is done being routed
        // and the onWhen predicate is true
        .to("log:sync")
        .to("mock:sync")
    // must use end to denote the end of the onCompletion route
    .end()
    // here the original route continues
    .to("log:original")
    .to("mock:result");
<route>
    <from uri="direct:start"/>
    <onCompletion>
        <onWhen>
            <simple>${body} contains 'Hello'</simple>
        </onWhen>
        <to uri="log:sync"/>
        <to uri="mock:sync"/>
    </onCompletion>
    <to uri="log:original"/>
    <to uri="mock:result"/>
</route>
- routeConfiguration:
    id: myConfig
    onCompletion:
      - onCompletion:
          onWhen:
            simple: "${body} contains 'Hello'"
          steps:
            - to:
                uri: log:sync
            - to:
                uri: mock:sync
- route:
    routeConfigurationId: myConfig
    from:
      uri: direct:start
      steps:
        - to:
            uri: log:original
        - to:
            uri: mock:result

使用 onCompletion(可选择是否启用线程池)

若要使用线程池,可以设置 executorService,或将 parallelProcessing 设置为 true。

  • Java
  • XML
  • YAML
onCompletion().parallelProcessing()
    .to("mock:before")
    .delay(1000)
    .setBody(simple("OnComplete:${body}"));
<onCompletion parallelProcessing="true">
  <to uri="mock:before"/>
  <delay><constant>1000</constant></delay>
  <setBody><simple>OnComplete:${body}</simple></setBody>
</onCompletion>

你还可以通过 executorService 选项指定要使用的特定线程池。

<onCompletion executorService="myThreadPool">
  <to uri="mock:before"/>
  <delay><constant>1000</constant></delay>
  <setBody><simple>OnComplete:${body}</simple></setBody>
</onCompletion>
- onCompletion:
    parallelProcessing: "true"
    steps:
      - to:
          uri: mock:before
      - delay:
          expression:
            constant:
              expression: 1000
      - setBody:
          expression:
            simple:
              expression: "OnComplete:${body}"

你还可以通过 executorService 选项指定要使用的特定线程池。

- onCompletion:
    executorService: myThreadPool
    steps:
      - to:
          uri: mock:before
      - delay:
          expression:
            constant:
              expression: 1000
      - setBody:
          expression:
            simple:
              expression: "OnComplete:${body}"

OnCompletion 消费者模式

OnCompletion 支持两种影响路由消费者的模式:

  • AfterConsumer - 默认模式,在消费者完成之后运行
  • BeforeConsumer - 在消费者完成之前运行,即在消费者将响应写回调用方之前运行

AfterConsumer 模式是默认模式,其行为与早期 Camel 版本相同。

新的 BeforeConsumer 模式用于在消费者将响应写回调用方之前(如果处于 InOut 模式)运行 onCompletion。这样 onCompletion 就可以修改 Exchange,例如添加特殊的报头,或者将 Exchange 记录为响应日志等。

例如,要始终添加一个「created by」报头,可以使用 modeBeforeConsumer(),如下所示:

  • Java
  • XML
  • YAML
.onCompletion().modeBeforeConsumer()
    .setHeader("createdBy", constant("Someone"))
.end()
<onCompletion mode="BeforeConsumer">
  <setHeader name="createdBy">
    <constant>Someone</constant>
  </setHeader>
</onCompletion>
- onCompletion:
    mode: BeforeConsumer
    steps:
      - setHeader:
          name: createdBy
          expression:
            constant:
              expression: Someone

评论

登录后参与评论

正在加载评论…