架构

虚拟线程

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

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

Apache Camel 中的虚拟线程

本指南介绍如何在 Apache Camel 中使用虚拟线程(Project Loom),以提升 I/O 密集型集成工作负载的性能。

简介

什么是虚拟线程?

虚拟线程于 JDK 19 中作为预览功能引入,并在 JDK 21 中正式定稿(JEP 444)。它们是由 JVM 而非操作系统管理的轻量级线程。借助虚拟线程,既能以熟悉的"每请求一个线程"的风格编写并发代码,又能获得异步编程的可扩展性。

关键特性

方面平台线程虚拟线程
管理方操作系统JVM
内存占用栈约 1 MB约 1 KB(按需增长)
创建开销昂贵(系统调用)廉价(对象分配)
最大实用数量数千条数百万条
阻塞行为阻塞操作系统线程挂起,释放承载线程

虚拟线程对集成的重要性

集成工作负载通常是 I/O 密集型 的——等待 HTTP 响应、数据库查询、消息代理确认或文件操作。使用平台线程时,每个阻塞操作都会占用一个昂贵的操作系统线程。而使用虚拟线程:

  • I/O 等待不会浪费资源——当虚拟线程因 I/O 阻塞时,它会"挂起",其承载线程可以运行其他虚拟线程
  • 大规模并发成为可能——无需担心线程池耗尽,即可处理数千个并发请求
  • 编程模型简单——编写直白的阻塞式代码,而不必编写复杂的响应式链

环境要求

  • JDK 21+ 用于虚拟线程
  • JDK 25+ 用于 ScopedValue 优化(可选,在上下文传播方面提供更好的性能)

在 Camel 中启用虚拟线程

在 Apache Camel 中,虚拟线程是 可选启用 的。启用后,Camel 的线程池工厂会自动为兼容的操作创建虚拟线程,而不是平台线程。

全局配置

Camel Main 配置属性

推荐通过 Camel Main 配置属性来启用虚拟线程。这种方式适用于 camel-main、camel-jbang 和 camel-spring-boot:

camel.main.virtualThreadsEnabled=true

该属性在引导过程的早期、线程池创建之前读取,因此会全局生效。

系统属性

另外,你也可以使用 JVM 系统属性:

java -Dcamel.threads.virtual.enabled=true -jar myapp.jar

编程方式配置

对于自定义配置,线程类型会在首次访问时根据系统属性延迟确定。Camel 的 ThreadType.current() 返回 PLATFORM 或 VIRTUAL。

启用后的变化

启用虚拟线程后,Camel 的 DefaultThreadPoolFactory(JDK 21+ 变体)的行为会发生变化:

线程池类型平台模式虚拟模式
newCachedThreadPool()Executors.newCachedThreadPool()Executors.newThreadPerTaskExecutor()
newThreadPool() (maxQueueSize > 0)ThreadPoolExecutor 与有界队列Executors.newThreadPerTaskExecutor() 并用基于信号量的并发限制进行包装
newThreadPool() (maxQueueSize ≤ 0)ThreadPoolExecutor 与 SynchronousQueueExecutors.newThreadPerTaskExecutor()(无界)
newScheduledThreadPool()ScheduledThreadPoolExecutorExecutors.newScheduledThreadPool(0, factory)

当 maxQueueSize 设置为正值时,虚拟线程执行器会被一个信号量包装,该信号量将并发上限强制固定为 maxQueueSize。与 ThreadPoolExecutor 中池线程和排队任务彼此独立不同,所有获得许可的任务都会立即在虚拟线程上执行。rejectedPolicy 控制达到并发限制时的行为——详情参见拒绝策略。

单线程执行器和计划任务仍然使用平台线程,因为虚拟线程针对的是并发 I/O 密集型工作,而非计划性或顺序性任务。

支持虚拟线程的组件

Camel 组件会因其架构不同而以不同方式受益于虚拟线程。

自动支持(基于线程池)

这些组件使用 Camel 的 ExecutorServiceManager,启用虚拟线程后可自动受益:

组件受益方式
SEDA / VM消费者线程变为虚拟线程;设置 virtualThreadPerTask=true 后,每条消息都会获得自己的虚拟线程
Direct-VM跨上下文调用使用虚拟线程进行异步处理
Threads DSL.threads() EIP 使用虚拟线程池
Async Processors使用 AsyncProcessor 并配合线程池的组件

HTTP 服务器组件

HTTP 服务器组件可以配置为使用虚拟线程来处理请求:

Jetty

Jetty 12+ 通过 VirtualThreadPool 支持虚拟线程。配置自定义线程池:

仅 Java:以编程方式配置 Jetty 的 VirtualThreadPool

import org.eclipse.jetty.util.thread.VirtualThreadPool;

JettyHttpComponent jetty = context.getComponent("jetty", JettyHttpComponent.class);

// Create Jetty's VirtualThreadPool for request handling
VirtualThreadPool virtualThreadPool = new VirtualThreadPool();
virtualThreadPool.setName("CamelJettyVirtual");
jetty.setThreadPool(virtualThreadPool);

或在 Spring 配置中:

<bean id="jettyThreadPool" class="org.eclipse.jetty.util.thread.VirtualThreadPool">
    <property name="name" value="CamelJettyVirtual"/>
</bean>

<bean id="jetty" class="org.apache.camel.component.jetty.JettyHttpComponent">
    <property name="threadPool" ref="jettyThreadPool"/>
</bean>

Platform HTTP (Vert.x)

camel-platform-http-vertx 组件使用 Vert.x 的事件循环模型。虚拟线程不适用于此场景,但你可以将阻塞工作卸载到其他线程执行:

  • Java
  • XML
  • YAML
from("platform-http:/api/orders")
    .threads()  // Offload to virtual thread pool
    .to("jpa:Order");  // Blocking JPA operation
<route>
    <from uri="platform-http:/api/orders"/>
    <threads/>
    <to uri="jpa:Order"/>
</route>
- route:
    from:
      uri: platform-http:/api/orders
      steps:
        - threads: {}
        - to:
            uri: jpa:Order

Undertow

Undertow 可以通过 XNIO worker 配置使用虚拟线程。有关 JDK 21+ 的虚拟线程支持,请参阅 Undertow 文档。

消息组件

组件虚拟线程使用方式
Kafka消费者线程池在高并发场景下可受益于虚拟线程
JMS会话处理和消息监听器可以使用虚拟线程池
AMQP连接处理可受益于虚拟线程

数据库组件

虚拟线程在阻塞性数据库操作中表现出色:

  • Java
  • XML
  • YAML
// With virtual threads, these blocking calls don't waste platform threads
from("seda:process?virtualThreadPerTask=true&concurrentConsumers=500")
    .to("jpa:Order")           // Blocking JDBC under the hood
    .to("sql:SELECT * FROM inventory WHERE id = :#${body.itemId}")
    .to("mongodb:orders");
<route>
    <from uri="seda:process?virtualThreadPerTask=true&amp;concurrentConsumers=500"/>
    <to uri="jpa:Order"/>
    <to uri="sql:SELECT * FROM inventory WHERE id = :#${body.itemId}"/>
    <to uri="mongodb:orders"/>
</route>
- route:
    from:
      uri: seda:process
      parameters:
        virtualThreadPerTask: true
        concurrentConsumers: 500
      steps:
        - to:
            uri: jpa:Order
        - to:
            uri: "sql:SELECT * FROM inventory WHERE id = :#${body.itemId}"
        - to:
            uri: mongodb:orders

SEDA 深入剖析:两种执行模型

Apache Camel 中的 SEDA(分阶段事件驱动架构)组件在路由之间提供异步的内存内消息传递。随着虚拟线程的引入,SEDA 现在支持两种不同的执行模型,每种模型都针对不同的场景进行了优化。

传统模型:固定消费者线程池

默认的 SEDA 消费者模型使用一组长期运行的固定消费者线程池,持续轮询队列以获取消息。

工作原理

  1. 消费者启动时,会创建 concurrentConsumers 个线程(默认值:1)
  2. 每个线程运行在一个无限循环中,以可配置的超时时间轮询队列
  3. 当消息到达时,线程处理该消息,然后再次进行轮询
  4. 线程会在多条消息之间被复用

配置

  • Java
  • XML
  • YAML
from("seda:orders?concurrentConsumers=10")
    .process(this::processOrder)
    .to("direct:fulfillment");
<route>
    <from uri="seda:orders?concurrentConsumers=10"/>
    <process ref="processOrder"/>
    <to uri="direct:fulfillment"/>
</route>
- route:
    from:
      uri: seda:orders
      parameters:
        concurrentConsumers: 10
      steps:
        - process:
            ref: processOrder
        - to:
            uri: direct:fulfillment

最佳适用场景

  • 受 CPU 限制的处理场景,其中线程创建开销不可忽视
  • 吞吐量可预测、稳定的场景
  • 需要精确控制线程池大小的场景
  • 平台线程(JDK < 21 或未启用虚拟线程)

每个任务一个虚拟线程模型

virtualThreadPerTask 模式采用了一种根本不同的方法:为每条消息创建一个新线程。

工作原理

  1. 单个协调器线程轮询队列
  2. 对于每条消息,向缓存线程池提交一个新任务
  3. 当启用虚拟线程时,使用 Executors.newThreadPerTaskExecutor()
  4. 每条消息都会获得自己的轻量级虚拟线程
  5. concurrentConsumers 选项变为并发上限(默认值:0 = 不限)

配置

  • Java
  • XML
  • YAML
from("seda:orders?virtualThreadPerTask=true&concurrentConsumers=100")
    .process(this::processOrder)  // I/O-bound operation
    .to("direct:fulfillment");
<route>
    <from uri="seda:orders?virtualThreadPerTask=true&amp;concurrentConsumers=100"/>
    <process ref="processOrder"/>
    <to uri="direct:fulfillment"/>
</route>
- route:
    from:
      uri: seda:orders
      parameters:
        virtualThreadPerTask: true
        concurrentConsumers: 100
      steps:
        - process:
            ref: processOrder
        - to:
            uri: direct:fulfillment

最佳适用场景

  • I/O 密集型工作负载(数据库调用、HTTP 请求、文件操作)
  • 吞吐量波动剧烈、流量突发的场景
  • 需要大规模并发的场景(数千个并发消息)
  • 虚拟线程(JDK 21+,且需设置 camel.threads.virtual.enabled=true)

架构对比

方面传统方式(固定线程池)每任务一个虚拟线程
线程创建启动时创建一次每条消息创建一次
线程数量固定(concurrentConsumers)动态(受上限约束)
队列轮询所有线程共同轮询单个协调者进行轮询
消息分发在轮询线程中直接处理提交给任务执行器
最适合CPU 密集型、平台线程I/O 密集型、虚拟线程
内存开销较高(平台线程约 1MB)较低(虚拟线程约 1KB)

可视化对比

flowchart TB
    subgraph traditional["Traditional Model (Fixed Pool)"]
        direction TB
        Q1[("SEDA Queue")]
        C1["Consumer Thread 1"]
        C2["Consumer Thread 2"]
        C3["Consumer Thread N"]
        P1["Process Message"]

        Q1 -->|"poll()"| C1
        Q1 -->|"poll()"| C2
        Q1 -->|"poll()"| C3
        C1 --> P1
        C2 --> P1
        C3 --> P1
    end

    subgraph virtual["Virtual Thread Per Task Model"]
        direction TB
        Q2[("SEDA Queue")]
        COORD["Coordinator Thread"]
        SEM{{"Semaphore (concurrency limit)"}}
        VT1["Virtual Thread 1"]
        VT2["Virtual Thread 2"]
        VTN["Virtual Thread N"]
        P2["Process Message"]

        Q2 -->|"poll()"| COORD
        COORD -->|"acquire"| SEM
        SEM -->|"spawn"| VT1
        SEM -->|"spawn"| VT2
        SEM -->|"spawn"| VTN
        VT1 --> P2
        VT2 --> P2
        VTN --> P2
    end

启用虚拟线程

要在 Camel 中使用虚拟线程,需要 JDK 21 或更高版本,并且必须通过配置启用它们:

应用属性

camel.threads.virtual.enabled=true

系统属性

java -Dcamel.threads.virtual.enabled=true -jar myapp.jar

启用后,Camel 的 DefaultThreadPoolFactory 会自动为缓存线程池使用 Executors.newThreadPerTaskExecutor(),从而创建虚拟线程而非平台线程。

背压与流量控制

在高并发场景下使用虚拟线程时,适当的背压机制至关重要,可避免给下游系统造成过大压力。SEDA 提供了多层背压控制。

第一层:基于队列的背压(生产者端)

SEDA 队列本身充当一个大小可配置的缓冲区:

  • Java
  • XML
  • YAML
// Queue holds up to 10,000 messages
from("seda:orders?size=10000")
<from uri="seda:orders?size=10000"/>
from:
  uri: seda:orders
  parameters:
    size: 10000

当队列已满时,生产者可以配置为:

选项行为使用场景
blockWhenFull=true生产者阻塞,直到有可用空间可以等待的同步调用方
blockWhenFull=true&offerTimeout=5000最多阻塞 5 秒,然后失败基于超时的流控
discardWhenFull=true静默丢弃消息即发即忘,允许丢失的场景
(默认)抛出 IllegalStateException快速失败,由调用方处理重试

阻塞与超时的示例:

  • Java
  • XML
  • YAML
// Producer blocks up to 10 seconds when queue is full
from("direct:incoming")
    .to("seda:processing?size=5000&blockWhenFull=true&offerTimeout=10000");
<route>
    <from uri="direct:incoming"/>
    <to uri="seda:processing?size=5000&amp;blockWhenFull=true&amp;offerTimeout=10000"/>
</route>
- route:
    from:
      uri: direct:incoming
      steps:
        - to:
            uri: seda:processing
            parameters:
              size: 5000
              blockWhenFull: true
              offerTimeout: 10000

第 2 层:并发限制(消费者端)

在 virtualThreadPerTask 模式下,concurrentConsumers 参数用于控制最大的并发处理任务数:

  • Java
  • XML
  • YAML
// Max 200 concurrent virtual threads processing messages
from("seda:orders?virtualThreadPerTask=true&concurrentConsumers=200")
    .to("http://downstream-service/api");
<route>
    <from uri="seda:orders?virtualThreadPerTask=true&amp;concurrentConsumers=200"/>
    <to uri="http://downstream-service/api"/>
</route>
- route:
    from:
      uri: seda:orders
      parameters:
        virtualThreadPerTask: true
        concurrentConsumers: 200
      steps:
        - to:
            uri: http://downstream-service/api

它在内部使用 Semaphore 来控制消息的分发,即使队列中积压了数千条消息,也能确保不会给下游服务造成过载。

第 3 层:组合策略

对于健壮的生产系统,建议将两者结合使用:

  • Java
  • XML
  • YAML
// Producer side: buffer up to 10,000, block if full (with timeout)
from("rest:post:/orders")
    .to("seda:order-queue?size=10000&blockWhenFull=true&offerTimeout=30000");

// Consumer side: process with virtual threads, max 500 concurrent
from("seda:order-queue?virtualThreadPerTask=true&concurrentConsumers=500")
    .to("http://inventory-service/check")
    .to("http://payment-service/process")
    .to("jpa:Order");
<!-- Producer side: buffer up to 10,000, block if full (with timeout) -->
<route>
    <from uri="rest:post:/orders"/>
    <to uri="seda:order-queue?size=10000&amp;blockWhenFull=true&amp;offerTimeout=30000"/>
</route>

<!-- Consumer side: process with virtual threads, max 500 concurrent -->
<route>
    <from uri="seda:order-queue?virtualThreadPerTask=true&amp;concurrentConsumers=500"/>
    <to uri="http://inventory-service/check"/>
    <to uri="http://payment-service/process"/>
    <to uri="jpa:Order"/>
</route>
# Producer side: buffer up to 10,000, block if full (with timeout)
- route:
    from:
      uri: rest:post:/orders
      steps:
        - to:
            uri: seda:order-queue
            parameters:
              size: 10000
              blockWhenFull: true
              offerTimeout: 30000

# Consumer side: process with virtual threads, max 500 concurrent
- route:
    from:
      uri: seda:order-queue
      parameters:
        virtualThreadPerTask: true
        concurrentConsumers: 500
      steps:
        - to:
            uri: http://inventory-service/check
        - to:
            uri: http://payment-service/process
        - to:
            uri: jpa:Order

该配置:

  • 在内存中最多缓冲 10,000 个订单
  • 缓冲区满时,最多阻塞 REST 调用方 30 秒
  • 最多以 500 个并发虚拟线程进行处理
  • 保护下游 HTTP 服务免于过载

背压对比

机制控制内容位置
size队列容量(消息缓冲区)生产者与消费者之间
blockWhenFull / offerTimeout生产者阻塞行为生产者一侧
concurrentConsumers(传统方式)固定线程池大小消费者一侧
concurrentConsumers(virtualThreadPerTask)最大并发任务数(信号量)消费者一侧

示例:高吞吐量订单处理

仅 Java 实现:包含 REST 与 SEDA 虚拟线程处理的 RouteBuilder 类

public class OrderProcessingRoute extends RouteBuilder {
    @Override
    public void configure() {
        // Receive orders via REST, queue them for async processing
        // Block callers if queue is full (with 30s timeout)
        rest("/orders")
            .post()
            .to("seda:incoming-orders?size=10000&blockWhenFull=true&offerTimeout=30000");

        // Process with virtual threads - each order gets its own thread
        // Limit to 500 concurrent to protect downstream services
        from("seda:incoming-orders?virtualThreadPerTask=true&concurrentConsumers=500")
            .routeId("order-processor")
            .log("Processing order ${body.orderId} on ${threadName}")
            .to("http://inventory-service/check")      // I/O - virtual thread parks
            .to("http://payment-service/process")      // I/O - virtual thread parks
            .to("jpa:Order")                           // I/O - virtual thread parks
            .to("direct:send-confirmation");
    }
}

性能特征

对于虚拟线程和 I/O 密集型工作负载,你可以期待:

  • 更高的吞吐量:虚拟线程在 I/O 等待期间不会阻塞操作系统线程
  • 更好的资源利用率:以极低的内存占用支持数千个并发操作
  • 负载下的更低延迟:不会出现线程池耗尽或排队延迟
  • 更简单的扩展方式:只需提高并发限制,无需调优线程池

基准测试

运行附带的负载测试以比较不同模式:

# Platform threads, fixed pool
mvn test -Dtest=VirtualThreadsLoadTest -pl core/camel-core

# Virtual threads, fixed pool
mvn test -Dtest=VirtualThreadsLoadTest -pl core/camel-core \
    -Dcamel.threads.virtual.enabled=true

# Virtual threads, thread-per-task (optimal)
mvn test -Dtest=VirtualThreadsLoadTest -pl core/camel-core \
    -Dcamel.threads.virtual.enabled=true \
    -Dloadtest.virtualThreadPerTask=true

Context Propagation with ContextValue

虚拟线程面临的一个挑战是上下文传播——即在整个调用链中传递上下文数据(如事务 ID、租户信息或用户凭据)。传统的 ThreadLocal 可以使用,但在虚拟线程场景下存在局限性。

ThreadLocal 的问题

ThreadLocal 在虚拟线程环境中存在以下问题:

  • 内存开销:每个虚拟线程都需要自己的一份副本
  • 继承复杂性:值必须显式继承给子线程
  • 无自动清理:如果值未被移除,存在泄漏风险
  • 没有作用域:值会一直存在,直到被显式移除

ContextValue 简介

Apache Camel 提供了 ContextValue 抽象,它会根据 JDK 版本和配置自动选择最优实现:

JDK 版本是否启用虚拟线程实现方式
JDK 17-24不适用ThreadLocal
JDK 21-24是ThreadLocal(ScopedValue 尚未稳定)
JDK 25+是ScopedValue
JDK 25+否ThreadLocal

ScopedValue 的优势(JDK 25+)

JEP 487:Scoped Values 提供了以下特性:

  • 不可变性:值在作用域内无法被修改(更安全)
  • 自动继承:子虚拟线程会自动继承值
  • 自动清理:离开作用域时值会被解除绑定(无泄漏)
  • 更优性能:针对结构化并发模型进行了优化

使用 ContextValue

基本用法

仅限 Java:用于作用域上下文传播的 ContextValue API

import org.apache.camel.util.concurrent.ContextValue;

// Create a context value (picks ScopedValue or ThreadLocal automatically)
private static final ContextValue<String> TENANT_ID = ContextValue.newInstance("tenantId");

// Bind a value for a scope
ContextValue.where(TENANT_ID, "acme-corp", () -> {
    // Code here can access TENANT_ID.get()
    processRequest();
    return result;
});

// Inside processRequest(), on any thread in the scope:
public void processRequest() {
    String tenant = TENANT_ID.get();  // Returns "acme-corp"
    // ... process with tenant context
}

何时使用 ThreadLocal 与 ContextValue

仅适用于 Java:在 ContextValue 工厂方法之间进行选择

// Use ContextValue.newInstance() for READ-ONLY context passing
private static final ContextValue<RequestContext> REQUEST_CTX = ContextValue.newInstance("requestCtx");

// Use ContextValue.newThreadLocal() when you need MUTABLE state
private static final ContextValue<Counter> COUNTER = ContextValue.newThreadLocal("counter", Counter::new);

与 Camel 内部机制集成

Camel 在内部出于多种目的使用 ContextValue:

仅限 Java:Camel 内部处理器创建过程中对 ContextValue 的使用

// Example: Passing context during route creation
private static final ContextValue<ProcessorDefinition<?>> CREATE_PROCESSOR
    = ContextValue.newInstance("CreateProcessor");

// When creating processors, bind the context
ContextValue.where(CREATE_PROCESSOR, this, () -> {
    return createOutputsProcessor(routeContext);
});

// Child code can access the current processor being created
ProcessorDefinition<?> current = CREATE_PROCESSOR.orElse(null);

从 ThreadLocal 迁移

如果你现有的代码使用了 ThreadLocal,迁移起来很简单:

仅 Java:从 ThreadLocal 迁移到 ContextValue

// Before: ThreadLocal
private static final ThreadLocal<User> CURRENT_USER = new ThreadLocal<>();

public void handleRequest(User user) {
    CURRENT_USER.set(user);
    try {
        processRequest();
    } finally {
        CURRENT_USER.remove();
    }
}

// After: ContextValue
private static final ContextValue<User> CURRENT_USER = ContextValue.newInstance("currentUser");

public void handleRequest(User user) {
    ContextValue.where(CURRENT_USER, user, this::processRequest);
}

ContextValue 版本更简洁,并且会自动处理清理工作。

最佳实践与性能注意事项

何时使用虚拟线程

适合 ✓不适合 ✗
HTTP 客户端调用CPU 密集型计算
数据库查询(JDBC)没有 I/O 的紧循环
文件 I/O 操作实时/低延迟系统
消息代理操作阻塞的原生代码(JNI)
调用外部服务长时间持有锁的代码

配置指南

从保守配置开始

# Start with virtual threads disabled, benchmark, then enable
camel.threads.virtual.enabled=false

# When enabling, test thoroughly
camel.threads.virtual.enabled=true

SEDA 调优

  • Java
  • XML
  • YAML
// For I/O-bound: use virtualThreadPerTask with high concurrency limit
from("seda:io-bound?virtualThreadPerTask=true&concurrentConsumers=1000")

// For CPU-bound: stick with traditional model, tune pool size
from("seda:cpu-bound?concurrentConsumers=4")  // ~number of CPU cores
<!-- For I/O-bound: use virtualThreadPerTask with high concurrency limit -->
<from uri="seda:io-bound?virtualThreadPerTask=true&amp;concurrentConsumers=1000"/>

<!-- For CPU-bound: stick with traditional model, tune pool size -->
<from uri="seda:cpu-bound?concurrentConsumers=4"/>
# For I/O-bound: use virtualThreadPerTask with high concurrency limit
from:
  uri: seda:io-bound
  parameters:
    virtualThreadPerTask: true
    concurrentConsumers: 1000

# For CPU-bound: stick with traditional model, tune pool size
from:
  uri: seda:cpu-bound
  parameters:
    concurrentConsumers: 4

避免固定

虚拟线程在以下情况下会被"固定"到承载线程(carrier thread)上:

  • 在 synchronized 代码块内部
  • 在进行本地方法(native method)调用期间

优先使用 ReentrantLock 而不是 synchronized:

仅 Java:使用 ReentrantLock 代替 synchronized 以避免固定

// Avoid: can pin virtual thread
synchronized (lock) {
    doBlockingOperation();
}

// Prefer: virtual thread can unmount
lock.lock();
try {
    doBlockingOperation();
} finally {
    lock.unlock();
}

监控与调试

线程名称

Camel 创建的虚拟线程具有描述性的名称:

VirtualThread[#123]/Camel (camel-1) thread #5 - seda://orders

JFR 事件

JDK Flight Recorder(JDK 飞行记录器)会捕获虚拟线程事件:

# Record virtual thread events
java -XX:StartFlightRecording=filename=recording.jfr,settings=default \
     -Dcamel.threads.virtual.enabled=true \
     -jar myapp.jar

检测线程钉住(Pinning)

# Log when virtual threads pin (JDK 21+)
java -Djdk.tracePinnedThreads=short \
     -Dcamel.threads.virtual.enabled=true \
     -jar myapp.jar

完整示例

示例 1:高并发 REST API

纯 Java 实现:用于高并发 REST API 的 RouteBuilder 类,基于虚拟线程

public class RestApiRoute extends RouteBuilder {
    @Override
    public void configure() {
        // REST endpoint receives requests
        rest("/api")
            .post("/orders")
            .to("seda:process-order");

        // Process with virtual threads - handle 1000s of concurrent requests
        from("seda:process-order?virtualThreadPerTask=true&concurrentConsumers=2000")
            .routeId("order-processor")
            // Each step may block on I/O - virtual threads park efficiently
            .to("http://inventory-service/reserve")
            .to("http://payment-service/charge")
            .to("jpa:Order?persistenceUnit=orders")
            .to("kafka:order-events");
    }
}

示例 2:使用虚拟线程进行并行增强

仅限 Java:RouteBuilder 类,通过编程方式创建执行器服务以实现并行增强

public class ParallelEnrichmentRoute extends RouteBuilder {
    @Override
    public void configure() {
        from("direct:enrich")
            .multicast()
                .parallelProcessing()
                .executorService(virtualThreadExecutor())  // Use virtual threads
                .to("direct:enrichFromUserService",
                    "direct:enrichFromOrderHistory",
                    "direct:enrichFromRecommendations")
            .end()
            .to("direct:aggregate");
    }

    private ExecutorService virtualThreadExecutor() {
        return getCamelContext()
            .getExecutorServiceManager()
            .newCachedThreadPool(this, "enrichment");
        // When camel.threads.virtual.enabled=true, this returns a virtual thread executor
    }
}

示例 3:路由内的上下文传播

仅限 Java:演示如何利用交换属性进行上下文传播的 RouteBuilder 类

public class TenantAwareRoute extends RouteBuilder {

    private static final ContextValue<String> TENANT_ID = ContextValue.newInstance("tenantId");

    @Override
    public void configure() {
        // ContextValue is scoped to the current thread - it works within a single
        // route or call chain, not across asynchronous boundaries like SEDA queues.
        // For cross-route context, use exchange properties instead.
        from("platform-http:/api/{tenant}/orders")
            .process(exchange -> {
                String tenant = exchange.getMessage().getHeader("tenant", String.class);
                exchange.setProperty("tenantId", tenant);
            })
            .to("seda:process");

        from("seda:process?virtualThreadPerTask=true&concurrentConsumers=500")
            .process(exchange -> {
                // Use exchange properties for context that crosses async boundaries
                String tenant = exchange.getProperty("tenantId", "default", String.class);
                log.info("Processing for tenant: {}", tenant);
            })
            .toD("jpa:Order?persistenceUnit=${exchangeProperty.tenantId}");
    }
}
ContextValue 的作用域仅限于当前线程(在 JDK 25+ 上则是 ScopedValue 的作用域)。它不会跨越 SEDA 队列等异步边界进行传播。对于需要跨越路由边界传递的数据,请使用交换属性或消息头。ContextValue 专为在同步调用链内部传播上下文而设计(例如在路由创建或处理器初始化期间)。

摘要

Apache Camel 中的虚拟线程提供了以下优势:

  • 简化的并发 - 编写阻塞代码,无需回调地狱
  • 更好的可扩展性 - 处理数千个并发 I/O 操作
  • 更低的资源消耗 - 轻量级线程占用更少的内存
  • 更高的吞吐量 - 高负载下不会出现线程池耗尽

快速上手:

  1. 升级到 JDK 21 及以上版本
  2. 在配置中添加 camel.threads.virtual.enabled=true
  3. 对于 SEDA 组件,针对 I/O 密集型工作负载可考虑设置 virtualThreadPerTask=true
  4. 使用 -Djdk.tracePinnedThreads=short 进行监控以发现问题

对于更高级的上下文传播需求(尤其是在 JDK 25+ 上),请使用 ContextValue 代替原始的 ThreadLocal。

评论

登录后参与评论

正在加载评论…