架构

线程模型

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

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

线程模型

Camel 中的线程模型基于一个可插拔的响应式路由引擎,以及来自 JDK 并发 API 的线程池。

本页重点介绍线程池。Camel 在多个地方使用了线程池,例如:

  • 若干 EIP 模式支持使用线程池来实现并发
  • 用于异步连接的 SEDA 组件
  • Camel 路由中的 Threads EIP
  • 多个组件开箱即用便使用了线程池,例如 JMS、Kafka、Netty 或 Jetty

线程池配置档案

默认情况下,当 Camel 需要创建线程池时,其池配置基于一个配置档案,即默认线程池配置档案。

默认配置档案开箱即用地预配置了以下设置:

选项默认值描述
poolSize10设置默认的核心线程数(线程池中至少保留的线程数)
keepAliveTime60设置空闲线程的默认存活时间(以秒为单位)
maxPoolSize20设置默认的最大线程池大小
maxQueueSize1000设置工作队列中默认的最大任务数。使用 -1 表示无界队列。
allowCoreThreadTimeOuttrue设置是否允许核心线程超时终止的默认值
rejectedPolicyCallerRuns设置线程池无法执行任务时的默认处理策略。共有三个选项:Abort、CallerRuns、Block。详情参见 拒绝策略。

这意味着,例如当你启用 parallelProcessing=true 使用 Multicast 时,它将基于上述配置创建一个线程池。

你可以定义任意多个线程池配置(profile),但默认配置只能有 一个。自定义线程池配置会继承自默认配置,也就是说,任何你没有显式定义的选项,都会回退并使用默认配置中的选项。

配置默认线程池配置

在 Spring XML 中,可以使用 threadPoolProfile 来配置线程池配置文件(thread pool profile),如下所示:

  • Java
  • Spring XML
  • Application Properties
ThreadPoolProfile profile = camelContext.getExecutorServiceManager().getDefaultThreadPoolProfile();
profile.setPoolSize(5);
profile.setMaxPoolSize(10);
<threadPoolProfile id="defaultThreadPoolProfile"
    defaultProfile="true"
    poolSize="5"
    maxPoolSize="10"/>

你也可以在 application.properties 中配置线程池。默认线程池可以按如下方式轻松配置:

## configure default thread pool profile
camel.threadpool.pool-size = 5
camel.threadpool.max-pool-size = 5

使用线程池配置文件

假设你想为 Camel 路由中的 Multicast EIP 模式使用自定义线程池配置文件,可以通过 executorService 属性实现,如下所示的 Spring XML:

  • Java
  • Spring XML
  • YAML

在 Java DSL 中,你可以使用 ThreadPoolProfileBuilder 来创建配置文件,然后注册该配置文件:

// setup thread pool profile
ThreadPoolProfileBuilder builder = new ThreadPoolProfileBuilder("fooProfile");
builder.poolSize(20).maxPoolSize(50).maxQueueSize(-1);

// add profile to camel
camelContext.getExecutorServiceManager().registerThreadPoolProfile(builder.build());

// camel routes can now refer to the thread pool via its id
from("direct:start")
  .multicast().aggregationStrategy("myStrategy").executorService("fooProfile").to("xxx");
<camelContext>

    <threadPoolProfile id="fooProfile"
                       poolSize="20" maxPoolSize="50" maxQueueSize="-1"/>

    <route>
       <multicast aggregationStrategy="myStrategy" executorService="fooProfile">
          ...
       </multicast>
    </route>
</camelContext>

要创建新的 profile,请使用 camel.threadpool.config[profileId] 语法,如下所示:

camel.threadpool.config[fooProfile].id = fooProfile
camel.threadpool.config[fooProfile].pool-size = 20
camel.threadpool.config[fooProfile].max-pool-size = 50
camel.threadpool.config[fooProfile].max-queue-size = -1

然后你可以在 Camel 路由中引用该线程池:

- route:
    from:
      uri: direct:start
      steps:
        - multicast:
            aggregationStrategy: myStrategy
            executorService: fooProfile
            steps:
              - to:
                  uri: log:hello

Camel 在运行时会到 Registry 中查找 id 为 fooProfile 的 ExecutorService。如果没有找到,它会回退并检查是否存在使用该 id 定义的 ThreadPoolProfile。在这个示例中存在这样一个 profile,因此该 profile 将作为基础设置来创建一个新的 ExecutorService,并将其返回给 Multicast EIP,以便在 Camel 路由中使用。

创建自定义线程池

你还可以在 Spring XML 中使用 <threadPool> 标签来创建 Java 线程池(ExecutorService)。请注意,任何你未显式定义的选项,Camel 都会回退使用默认线程池 profile。例如,如果你省略了 maxQueueSize 的设置,Camel 将回退使用默认线程池 profile 中的值,该值默认为 1000。

Spring XML

<camelContext>

    <threadPool id="myPool" threadName="myPool" poolSize="20" maxPoolSize="50"/>

    <route>
       <from uri="direct:start"/>
       <multicast aggregationStrategy="myStrategy" executorService="myPool">
          ...
       </multicast>
    </route>
</camelContext>

线程池与线程池配置文件有什么区别?

线程池配置文件(thread pool profile)是用于创建一个或多个实际线程池(即 ExecutorService)的模板。线程池则是具体的 ExecutorService 对象。

这样做的目的是通过配置文件,方便地使用相同设置创建多个线程池。

自定义线程名称

在 ExecutorServiceManager 上,可以使用 setThreadNamePattern 方法配置线程名称模式,该模式定义了线程池创建线程时所使用的线程名称。

默认模式为:

Camel (#camelId#) thread ##counter# - #name#

在该模式中,你可以使用以下占位符:

  • #camelId# - CamelContext 的名称
  • #counter# 一个唯一递增的计数器
  • #name# - 线程名称
  • #longName# - 长线程名称,可包含端点参数等
  • Java
  • Spring XML
  • Application Properties

在 Java 中,你可以在 ExecutorServiceManager 上设置该模式,如下所示:

CamelContext camel = ...
camel.getExecutorServiceManager().setThreadNamePattern("Riding the thread #counter#")

在 Spring XML 中,可以通过 threadNamePattern 属性来设置该模式,如下所示:

<camelContext threadNamePattern="Riding the thread #counter#">
  <route>
    <from uri="seda:start"/>
    <to uri="log:result"/>
    <to uri="mock:result"/>
  </route>
</camelContext>

使用 camel-main、Spring Boot 或 Quarkus,可以在 application.properties|yaml 文件中进行配置:

camel.main.thread-name-pattern = Riding the thread #counter#

关闭线程池

当 CamelContext 关闭时,Camel 创建的所有线程池都会被正常关闭,这可确保在支持热部署等特性的服务器环境中运行时,线程池不会发生泄漏。

ExecutorServiceManager 提供了用于优雅关闭和强制关闭线程池的 API。建议使用该 API 来创建和关闭线程池。

ExecutorServiceManager 的 shutdownGraceful(executorService) 方法会先执行优雅关闭,直到达到超时值为止。超时之后,它将执行强制关闭,同样使用超时值来等待操作完成。这意味着关闭线程池的等待时间最多为 2 × 超时值。

超时值默认为 10000 毫秒。如有需要,可以在 ExecutorServiceManager 上配置自定义值。在关闭过程中,Camel 会以 INFO 级别每隔 2 秒记录一次线程池关闭的进度。例如,如果关闭过程耗时较长,日志中就会体现出相应的活动。

与关闭线程池相关的 ExecutorServiceManager API 如下所示:

方法说明
shutdown将线程池标记为已关闭(类似于调用 ExecutorService.shutdown() 方法)。
shutdownNow强制线程池立即关闭(类似于调用 ExecutorService.shutdownNow() 方法)。
shutdownGraceful将线程池标记为已关闭,并优雅地关闭线程池,即等待任务完成。默认超时值为 10 秒,超时后将转为强制关闭,调用 shutdownNow 以使线程更快地关闭。
shutdownGraceful(timeout)与 shutdownGraceful 相同,但使用自定义的超时值
awaitTermination等待线程池优雅地终止(例如,等待其任务完成)。将一直等待,直到所有任务完成或超时。

JMX 管理

Camel 创建的所有线程池都受到管理,因此你可以在 JMX 的 threadpools 树下看到它们。

这需要通过将 camel-management JAR 包含在 classpath 中来启用 JMX。

组件开发者

如果你开发自己的 Camel 组件并且需要一个线程池,建议使用 ExecutorServiceStrategy/ExecutorServiceManager 来创建所需的线程池。

ExecutorServiceStrategy 与 ExecutorServiceManager

Camel 提供了一种可插拔的策略,用于挂接你自己的线程池提供程序。请参阅 org.apache.camel.spi.ExecutorServiceStrategy 接口,你需要实现该接口并将其挂接到 WorkManager 中。

若要接入自定义线程池提供程序,可以实现 ThreadPoolFactory 接口。该实现可在 ExecutorServiceManager 中进行设置。

拒绝策略

rejectedPolicy 选项用于控制当线程池无法接受新任务(即线程池及其工作队列已满)时的行为。可用的策略如下:

策略说明
CallerRuns任务在调用者的线程上运行。这提供了自然的背压机制——调用者被阻塞去执行实际工作,在完成之前无法提交更多任务。任务永远不会丢失。这是默认策略。
Abort任务被拒绝,并抛出 RejectedExecutionException。对于 HTTP API 或延迟敏感的系统,当快速失败比阻塞更可取时,可使用此策略。
Block调用者无限期阻塞,直到有可用容量。既不超时,也不拒绝。对于消息代理消费者和批处理工作负载,当丢失任务不可接受且延迟不那么关键时,可使用此策略。

使用平台线程时,当 ThreadPoolExecutor 的工作队列已满会应用这些策略。使用虚拟线程时,当并发信号量没有可用许可时会应用相同的策略(参见虚拟线程)。

虚拟线程

从 Java 21 开始,默认的 ThreadPoolFactory 可以构建使用虚拟线程而非平台线程的 ExecutorService 和 ScheduledExecutorService。

要启用虚拟线程,请将系统属性 camel.threads.virtual.enabled 设置为 true,并使用 Java 21 或更高版本运行 Camel。

请注意,即使启用了该功能,在某些用例中仍会使用平台线程,例如:当线程工厂被配置为创建非守护线程时(因为虚拟线程只能是守护线程),或者要构建的 ExecutorService 或 ScheduledExecutorService 不允许拥有多个线程时,又或者当 corePoolSize 设置为零且 maxQueueSize 设置为小于或等于 0 的值时。

虚拟线程的有界并发

当 maxQueueSize 设置为正值时,Camel 会用基于信号量的并发限制包装虚拟线程执行器。这样即使虚拟线程不使用传统的工作队列,也能确保 maxQueueSize 在背压场景下得到遵守。

与 ThreadPoolExecutor 中池线程和排队任务是不同概念不同,虚拟线程执行器采用的是扁平的并发上限:并发执行的最大任务数等于 maxQueueSize。线程池大小参数(poolSize、maxPoolSize)会被忽略,因为虚拟线程不会被池化。所有被允许的任务都会立即在虚拟线程上执行——不存在等待任务的队列。

rejectedPolicy 控制当达到并发上限时的行为:

  • CallerRuns(默认):调用方阻塞等待许可证,最长等待 keepAliveTime。如果超时,任务将在调用方的线程上执行。
  • Abort:调用方阻塞等待许可证,最长等待 keepAliveTime。如果超时,抛出 RejectedExecutionException。
  • Block:调用方无限期阻塞,直到获得许可证。

在等待许可证期间,调用线程会被阻塞。当调用方是虚拟线程时,这开销很小(承载线程会被释放)。当调用方是平台线程(例如 HTTP 服务器线程)时,被阻塞的线程就无法处理其他工作。

线程池大小参数(poolSize、maxPoolSize、keepAliveTime)不会控制虚拟线程的线程复用,因为虚拟线程创建成本低且从不池化。不过,在使用 CallerRuns 和 Abort 策略时,keepAliveTime 会被复用为获取许可证的超时时间。

评论

登录后参与评论

正在加载评论…