Ohhnews

分类导航

$ cd ..
DZone Java原文

使用Java虚拟线程与JMS监听器的实用指南

#java虚拟线程#jms#spring框架#并发处理#消息队列

使用 Java 虚拟线程扩展 JMS 监听器

事件驱动架构广泛用于企业系统,以解耦服务、吸收流量峰值并将工作移出请求路径。Java 消息服务(JMS)现已标准化为 Jakarta Messaging,在围绕 ActiveMQ、IBM MQ、Solace、TIBCO EMS 及类似代理构建的系统中仍然很常见。Java 21 虚拟线程为这些系统提供了另一种扩展选项。JMS 监听器通常在数据库、HTTP 服务、缓存或文件系统上等待的时间多于使用 CPU 的时间。将这种阻塞工作转移到虚拟线程可以减轻平台线程压力,而无需将应用程序强制改为反应式编程模型。然而,虚拟线程并不会让代理、数据库或下游服务变得无限。它们也不会改变确认、事务、重新投递或排序语义。安全的设计将虚拟线程与有界的 JMS 消费者并发、明确的资源限制、幂等性和生产度量相结合。本文解释虚拟线程改变了 Spring JMS 监听器的哪些方面,如何显式配置它们,以及如何避免将瓶颈从 JVM 转移到系统的其他部分。

传统的 JMS 监听器模型

典型的基于队列的流程将消息从代理通过 Spring 监听器容器移动到调用下游系统的处理器中。图 1 对比了该处理器工作在平台线程上的占用情况,以及当容器的消费者调用器任务使用虚拟线程时它的运行方式。

[LOADING...]

图 1. Spring JMS 监听器中平台线程与虚拟线程消费者调用器的对比。

容器管理 JMS 连接、会话、消费者、确认和监听器调用。处理器包含业务逻辑:

$ java
@JmsListener(
        destination = "orders.created",
        containerFactory = "jmsListenerContainerFactory"
)
public void handle(OrderCreatedEvent event) {
    Customer customer = customerClient.getCustomer(event.customerId());
    inventoryService.reserve(event.orderId(), customer);
    orderRepository.markAsProcessing(event.orderId());
}

这段代码易于阅读,但每个下游操作都可能阻塞。使用平台线程时,在等待查询或网络调用期间,操作系统后台线程会一直处于占用状态。当足够多的监听器线程被阻塞时,即使 CPU 未饱和,新消息也会等待。应用程序已经变成线程绑定型,而不是 CPU 绑定型。在虚拟线程出现之前,团队通常会增大监听器线程池、横向扩展更多服务实例,或者改用异步或响应式 API 重写流程。这些选项仍然有效,但每种都有代价。更大的平台线程池会占用更多内存并增加调度开销。更多实例会增加基础设施和运维工作。响应式代码可以高效扩展,但会改变库、控制流、调试和错误处理方式。

虚拟线程改变了什么

虚拟线程仍然是java.lang.Thread,但它由 JVM 调度,而不是永久绑定到一个操作系统线程。临时运行虚拟线程的平台线程被称为其载体(carrier)。当虚拟线程在受支持的 I/O 上阻塞时,JVM 可以将其从载体上卸载。载体随后可以自由运行另一个虚拟线程。这使得应用程序能够在支持大量并发阻塞操作的同时,保持直接、顺序的代码。如图 1 所示,等待受支持 I/O 的虚拟线程可以从其载体上卸载,从而使这些载体可以执行其他就绪工作。

当平台线程稀缺是限制因素时,虚拟线程可以提高吞吐量。它们不会使单个数据库调用或 HTTP 请求更快,也不会增加 CPU 容量。适合的候选包括由以下类型主导的处理器:

  • JDBC 调用
  • 阻塞式 REST 或 gRPC 客户端
  • 缓存查找
  • 文件或对象存储操作
  • 旧式同步 SDK
  • 跨下游系统的同步编排

不适合的候选包括由以下类型主导的处理器:

  • CPU 密集型转换
  • 加密或压缩
  • 图像或视频处理
  • 机器学习推理
  • 大规模内存聚合

JDK 的指导是每个任务创建一个虚拟线程,而不是对虚拟线程进行池化。有限资源应通过显式机制加以保护,例如信号量、速率限制器、连接池和框架并发设置。

改变设计的 JMS 细节

对于 Spring 的 DefaultMessageListenerContainer,监听器线程通常属于一个消费者调用器(consumer invoker)。该调用器拥有或复用 JMS SessionMessageConsumer,并且在其生命周期内可能会处理许多消息。因此,启用虚拟线程并不一定会为每条消息创建一个新的虚拟线程。它只是将容器的消费者任务放到虚拟线程上。这一区别很重要,因为提高并发度也会增加活动 JMS 消费者和会话的数量。这些代理端资源不像虚拟线程那样廉价。图 1 的右侧显式地建模了这一关系:一个配置好的消费者调用器任务运行在虚拟线程上,并可能在其生命周期内处理多条消息。这种架构仍然有用。当处理器在下游 I/O 上等待时,消费者可以从其载体上卸载。但监听器容器的并发度仍然是控制一次能处理多少条消息的主要手段。

显式配置 JMS 执行器

Spring Boot 可以通过 spring.threads.virtual.enabled=true 为多个由 Boot 管理的执行路径启用虚拟线程。不要仅凭这个属性就认为 JMS 监听器容器已使用虚拟线程。请显式配置 JMS 容器的执行器,并在运行时验证。图 2 将应用程序装配与运行时流程分开。启用虚拟线程的 TaskExecutor 与 JMS 监听器工厂之间的显式连接是关键步骤;容器的并发设置继续限制活动消费者和会话。

[LOADING...]

图 2. Spring JMS 虚拟线程显式装配与运行时消息流程。

以下示例使用 Java 21 或更高版本以及 Spring Framework 6.1 或更高版本。它向监听器容器工厂提供一个启用虚拟线程的 SimpleAsyncTaskExecutor

$ java
import java.util.concurrent.Executor;
import jakarta.jms.ConnectionFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.jms.config.DefaultJmsListenerContainerFactory;

@Configuration(proxyBeanMethods = false)
class JmsConfiguration {

    @Bean("jmsVirtualThreadExecutor")
    SimpleAsyncTaskExecutor jmsVirtualThreadExecutor() {
        SimpleAsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("jms-vt-");
        executor.setVirtualThreads(true);
        return executor;
    }

    @Bean
    DefaultJmsListenerContainerFactory jmsListenerContainerFactory(
            ConnectionFactory connectionFactory,
            @Qualifier("jmsVirtualThreadExecutor") Executor executor
    ) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        factory.setTaskExecutor(executor);
        // Example limits only. Derive these from load tests and
        // the safe capacity of the broker and downstream systems.
        factory.setConcurrency("10-100");
        // Prefer transactional JMS acknowledgment when redelivery
        // on listener failure is required.
        factory.setSessionTransacted(true);
        return factory;
    }
}

SimpleAsyncTaskExecutor.setVirtualThreads(true) 需要 Java 21。Spring Framework 6.2 还为那些直接构造监听器容器并使用其内部默认执行器的应用程序增加了 DefaultMessageListenerContainer.setVirtualThreads(true)。如果 Spring Boot 应用程序使用 Boot 的 DefaultJmsListenerContainerFactoryConfigurer,请在显式执行器、并发和事务覆盖之前应用它,以保留其他 Boot JMS 属性。虚拟线程是守护线程。在没有其他非守护线程维持 JVM 存活的非 Web 工作进程中,请使用 Spring Boot 的 spring.main.keep-alive=true 或等效的应用程序生命周期机制。不要依赖代理客户端产生的间接线程来保持进程运行。一个小的启动测试可以确认执行模式:

$ java
if (!Thread.currentThread().isVirtual()) {
    throw new IllegalStateException(
            "The JMS listener is not running on a virtual thread"
    );
}

请将其作为测试或临时诊断手段,而不是对每一条生产消息都执行。此外,当应用程序定义了多个容器工厂时,请确认当前使用的容器工厂。

根据实际容量限制并发

虚拟线程减少了线程稀缺,但并未消除资源稀缺。监听器仍可能受到以下因素的限制:

  • JMS 会话和消费者
  • 代理预取、消费者窗口或信用额度
  • 数据库连接
  • HTTP 客户端连接
  • 下游速率限制
  • 在途消息占用的内存
  • 事务锁
  • CPU

一个有用的初步估算来自利特尔定律(Little's Law):所需并发度 ≈ 目标吞吐量 × 平均处理时间 如果目标是每秒 200 条消息,平均处理器时间为 250 毫秒,则初始估计值为:

200 条消息/秒 × 0.25 秒 = 50 个并发处理器

该值只是一个起点。它必须受到每个依赖项安全容量的上限约束。如果每条消息持有一个数据库连接,而可用连接池容量为 30,那么将监听器并发度设置为 100 可能只会增加 70 个等待者。如果支付 API 允许 40 个并发请求,请使用信号量或速率限制器单独保护该调用。示例中的并发范围 10-100 意味着容器可以维持一个基线并扩展到最大值。它并不保证 100 是安全的,而且对于某些代理或工作负载,100 的上限可能过高。代理的流控设置也很重要。过度的预取(prefetch)可能会将大量积压从代理移入消费者,增加未确认消息的数量,并使恢复变得难以预测。保持足够的预取工作以喂饱消费者,但避免将预取当作无界应用队列来使用。

确认与事务必须刻意设计

虚拟线程不会改变消息投递保证。这对 Spring 的 DefaultMessageListenerContainer 尤为重要。在其默认的 AUTO_ACKNOWLEDGE 模式下,容器会在监听器执行之前进行确认,因此监听器异常不会引发重新投递。如果应用程序需要在处理器失败后回滚并重新投递,请使用事务性 JMS 会话或配置适当的外部事务管理器。本地 JMS 事务覆盖通过同一会话执行的 JMS 接收和 JMS 发送。它不会自动包含数据库事务。数据库提交可能成功,而 JMS 提交可能失败,从而导致消息被再次投递。有三种常见策略:

  1. 使用幂等处理器和本地事务。
  2. 使用 inbox/outbox 设计,使数据库效果可重复、对外发布可靠。
  3. 当需要 JMS 与另一个事务性资源之间进行原子协调,且其运维成本合理时,使用 JTA/XA。

图 3 展示了 inbox/outbox 生命周期,包括重复路径、独立的 JMS 确认边界、代理管理的重新投递和死信处理。

[LOADING...]

图 3. 幂等 JMS 处理、确认、重试和死信生命周期。

不要将数据库服务上的 @Transactional 视为 JMS 确认参与同一事务的证据。请核实哪个事务管理器处于活动状态,以及它协调哪些资源。

让消费者具备幂等性

重新投递可能发生在代理故障转移、事务回滚、应用程序重启、超时或两个资源提交之间发生故障之后。更高的并发度也更容易暴露重复检测中的竞态条件。inbox 表是一种常见解决方案。如图 3 所示,应用程序在同一数据库事务中原子性地插入消息 ID 并应用业务变更。重复键会走一条安全的空操作路径,而不是重复执行业务效果。数据库必须对消息 ID 强制唯一约束。仅靠单独的 exists() 检查是不够的,因为两个并发投递可能都观察到该行不存在。

$ java
@Transactional
public void process(OrderCreatedEvent event) {
    boolean firstDelivery = processedMessageRepository.tryInsert(event.messageId());
    if (!firstDelivery) {
        return;
    }
    orderService.apply(event);
}

tryInsert 应使用由唯一键保护的原子插入(insert-if-absent)操作,并在不提交单独事务的情况下报告重复。如果持久化提供程序将整个事务标记为仅回滚(rollback-only),请避免捕获通用的约束异常。如果业务更新失败,事务应同时回滚 inbox 插入和业务变更。外部副作用需要自己的幂等策略。例如,向支付 API 发送幂等键,或在调用无法参与本地事务的服务之前持久化操作状态。

保持事务和重试短小

避免在慢速外部服务重试数分钟时,一直持有数据库或 JMS 事务。风险模式是:开始事务、调用外部 API、等待并重试,然后才更新数据库并提交。这会持有锁、数据库连接、JMS 会话和未确认的消息。虚拟线程使等待线程更廉价,但并不会释放这些资源。更安全的设计(如图 3 所示)将业务更新和 outbox 记录作为本地意图提交,并通过 outbox 发布者异步继续处理。数据库更新和 outbox 插入发生在同一个本地事务中。一个单独的发布者发送待处理的 outbox 记录并将其标记为完成。如果在数据库提交后入站 JMS 消息被重新投递,inbox 键会防止业务更新和 outbox 插入被重复执行。长时间的重试延迟通常应通过代理的重新投递延迟、重试队列或调度器来处理。从载体线程的角度看,让虚拟线程休眠很廉价,但在延迟期间,监听器仍可能持有 JMS 消费者、会话、事务和消息。在重试之前应对错误进行分类:

失败类型典型应对
瞬时网络或依赖故障使用指数退避和抖动重试
速率限制遵守服务器的延迟并降低并发
无效消息模式发送到死信队列
缺少必要的业务数据进入死信或路由以进行修正
反复出现的未知故障在有限尝试次数后停止并告警

每个生产监听器都应定义最大重新投递次数、死信目标、重放流程,以及负责调查毒消息(poison message)的所有者。

不要随意将工作从监听器剥离

一个诱人的设计是让 JMS 监听器接收消息,将实际工作提交给另一个执行器,然后立即返回。这可以带来更多并行性,但也可能在工作完成之前确认消息。它还可能导致 JMS Session 跨越线程边界,而 Session 按约定是单线程的。事务上下文、错误传播和重新投递行为都可能会丢失。除非应用程序有意实现交接协议,否则应让监听器容器拥有处理器的执行权。安全的交接通常意味着在监听器返回之前持久化消息或命令,而不仅仅是将 Runnable 放入内存执行器。

在需要时保持有序性

更高的并发会改变排序行为。一旦队列有多个活动消费者,消息完成顺序可能与代理投递顺序不同。请显式地选择排序范围:

  • 对于严格的全局排序,将并发度保持为 1。
  • 按业务键进行分区或路由消息。
  • 对同一键的处理进行串行化。
  • 当事件可能乱序到达时,添加序号检查。
  • 设计状态转换以拒绝过期事件。

当消息彼此独立,或排序仅限于某个分区或业务键时,虚拟线程最容易采用。对于主题(topic),不要像对待队列那样增加消费者并发。根据订阅配置,额外的主题消费者可能会收到每条消息的额外副本。请审查代理和容器的持久订阅与共享订阅语义。

测试瓶颈,而不仅仅是线程数量

一个示例性的订单处理工作负载可能对每条消息执行一次数据库读取、两次 HTTP 调用、一次数据库更新和一次出站事件。在以下条件下比较平台线程和虚拟线程:

  • 相同的消息语料和负载分布
  • 相同的确认和事务设置
  • 相同的数据库和 HTTP 池限制
  • 相同的代理预取或信用额度
  • 相同的重试和死信策略
  • 受控的并发斜坡

不要只衡量吞吐量:

度量项揭示内容
队列深度和最早消息年龄积压和用户可见延迟
消费速率可持续吞吐量
处理器 p50、p95、p99 延迟正常和尾部行为
已调度和活动的 JMS 消费者实际容器并发度
平台线程和虚拟线程数量线程压力是否转移
载体 CPU 和线程固定事件调度器或兼容性问题
数据库池利用率和等待时间数据库饱和程度
HTTP 池利用率和超时出站连接压力
下游节流速率限制压力
重新投递和 DLQ 数量故障放大情况
堆和垃圾回收在途工作的成本

当系统以更低的平台线程压力维持所需吞吐量,并且没有增加超时、节流、重新投递或尾部延迟时,虚拟线程就是成功的。如果吞吐量上升而下游错误上升更快,系统并没有变得更健康,只是更高效地传递过载。

诊断线程固定与提供程序兼容性

在 Java 21 上,虚拟线程在执行某些 synchronized 或原生代码时阻塞,可能会固定(pin)其载体。偶发的短暂固定通常是无害的。频繁的长时间固定会降低可扩展性。使用 Java Flight Recorder 的 jdk.VirtualThreadPinned 事件,或使用以下参数运行负载测试:

-Djdk.tracePinnedThreads=full

请使用生产中实际使用的 JMS 提供程序、JDBC 驱动、HTTP 客户端、监控代理和安全库来进行测试。兼容性无法从合成的 Thread.sleep 基准测试推断出来。JDK 24 的 JEP 491 几乎消除了由 synchronized 方法和代码块引起的所有固定问题,但原生或外部函数交互以及第三方行为仍然值得测试。

决策矩阵

场景虚拟线程适配度
阻塞式 JDBC 调用
阻塞式 REST 或 gRPC 调用
旧式同步 SDK
高吞吐、I/O 密集的队列监听器强(需有界消费者)
CPU 密集型转换
严格的全局排序有限
下游容量小仅在严格限制下有用
较弱的确认或重试设计先修复投递语义
缺少可观测性先添加度量

生产清单

在为 JMS 监听器启用虚拟线程之前,请确认:

  • 应用程序运行在 Java 21 或更高版本。
  • JMS 执行器已显式配置并验证为虚拟。
  • 监听器并发度由测量到的下游容量限制。
  • 代理预取、消费者窗口或信用额度已调优。
  • 确认和事务行为已记录并测试。
  • 通过原子幂等机制防止重复处理。
  • 重试有界、有延迟且已分类。
  • 存在死信队列和重放流程。
  • 排序需求明确。
  • 负载测试使用真实驱动和代表性依赖。
  • 监控队列年龄、尾部延迟、池饱和、重新投递和线程固定事件。

结论

对于大部分时间都在等待阻塞 I/O 的 JMS 监听器来说,虚拟线程非常合适。它们让团队能够保留简单、命令式的 Java 代码,同时降低并发消息处理的平台线程成本。安全采用的模式不是“打开虚拟线程并移除限制”,而是:

  1. 将监听器容器的消费者任务放到虚拟线程上。
  2. 使用代理和下游容量限制消费者并发。
  3. 明确设计确认、事务和幂等性。
  4. 使用真实提供程序和依赖进行测试。
  5. 度量瓶颈移动到了哪里。

当这些控制措施到位后,虚拟线程可以在不要求响应式重写的情况下,使现有 JMS 应用程序现代化。它们让等待变得更廉价,但架构仍然需要决定系统能够安全接受多少工作。

参考资料

DZone 贡献者所表达的观点属于他们个人。