线程模型
线程模型
Camel 中的线程模型基于一个可插拔的响应式路由引擎,以及来自 JDK 并发 API 的线程池。
本页重点介绍线程池。Camel 在多个地方使用了线程池,例如:
- 若干 EIP 模式支持使用线程池来实现并发
- 用于异步连接的 SEDA 组件
- Camel 路由中的 Threads EIP
- 多个组件开箱即用便使用了线程池,例如 JMS、Kafka、Netty 或 Jetty
线程池配置档案
默认情况下,当 Camel 需要创建线程池时,其池配置基于一个配置档案,即默认线程池配置档案。
默认配置档案开箱即用地预配置了以下设置:
| 选项 | 默认值 | 描述 |
|---|---|---|
| poolSize | 10 | 设置默认的核心线程数(线程池中至少保留的线程数) |
| keepAliveTime | 60 | 设置空闲线程的默认存活时间(以秒为单位) |
| maxPoolSize | 20 | 设置默认的最大线程池大小 |
| maxQueueSize | 1000 | 设置工作队列中默认的最大任务数。使用 -1 表示无界队列。 |
| allowCoreThreadTimeOut | true | 设置是否允许核心线程超时终止的默认值 |
| rejectedPolicy | CallerRuns | 设置线程池无法执行任务时的默认处理策略。共有三个选项: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:helloCamel 在运行时会到 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 会被复用为获取许可证的超时时间。
评论
登录后参与评论
KnowForge