You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus中RestEasy端点并发执行阻塞校验任务的问题

Quarkus阻塞任务并发执行与上下文传播问题

问题背景

在Quarkus应用的RestEasy端点中,需发起一系列下游阻塞调用,期望通过并发执行提升效率,但当前实现里validateTrade(tradeRequest)始终串行执行,同时要保证JWT上下文能正确传播到工作线程。

当前实现代码

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) {
    // 将TradeRequest列表转换为Multi
    return Multi.createFrom().iterable(sendTrades)
            // 指定订阅处理的线程池
            .runSubscriptionOn(Infrastructure.getDefaultWorkerPool())
            // 将每个TradeRequest转换为Uni<Trade>,异步验证交易
            .onItem().transformToUniAndMerge(tradeRequest ->
                    Uni.createFrom().item(() -> validateTrade(tradeRequest))
            )
            // 将结果收集为列表
            .collect().asList()
            .await().atMost(Duration.ofMinutes(2));
}

已尝试使用Infrastructure.getDefaultWorkerPool()或自定义ManagedExecutor:

@Inject
@ManagedExecutorConfig(maxAsync = 5, maxQueued = 5)
@NamedInstance("tradeManagedExecutor")
ManagedExecutor managedExecutor;

尝试过的其他方案及结果

方案1:Uni组合+默认工作线程池

仍串行执行

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) {
        return Uni.combine().all()
                .unis(sendTrades.stream().map(trade -> Uni.createFrom().item(() -> validateTrade(trade))).toList())
                .combinedWith(listOfResponses -> listOfResponses.stream().map(Trade.class::cast).toList())
                .emitOn(Infrastructure.getDefaultWorkerPool())
                .await().atMost(Duration.ofMinutes(2));
    }

方案2:并行流

可并行执行但报错,提示请求上下文未激活

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) {
        return sendTrades.stream()
                .parallel()
                .map(this::validateTrade)
                .toList();
    }

错误信息:

Caused by: jakarta.enterprise.context.ContextNotActiveException: RequestScoped context was not active when trying to obtain a bean instance for a client proxy of CLASS bean [class=se.sebank.portfoliomgt.tradingservice.instrument.domain.InstrumentService, id=56969d55d00200c862f05c62808ce63ddb2d7514]
        - you can activate the request context for a specific method using the @ActivateRequestContext interceptor binding
        at io.quarkus.arc.impl.ClientProxies.getDelegate(ClientProxies.java:55)

方案3:CompletableFuture+自定义Executor

出现上下文传播问题,报空指针

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) throws ExecutionException, InterruptedException {
        List<CompletableFuture<Trade>> futureTrades = sendTrades.stream()
                .map(trade -> CompletableFuture.supplyAsync(() -> this.validateTrade(trade), managedExecutor))
                .toList();

        CompletableFuture<Void> allDoneFuture = CompletableFuture.allOf(futureTrades.toArray(new CompletableFuture[0]));

        return allDoneFuture.thenApply(v ->
                futureTrades.stream()
                        .map(CompletableFuture::join) // 获取每个CompletableFuture的结果
                        .toList()).get();
    }

错误信息:

Error invoking subclass method
        at org.jboss.resteasy.core.ExceptionHandler.handleApplicationException(ExceptionHandler.java:107)
        at org.jboss.resteasy.core.ExceptionHandler.handleException(ExceptionHandler.java:344)
        at org.jboss.resteasy.core.SynchronousDispatcher.writeException(SynchronousDispatcher.java:205)
        at org.jboss.resteasy.core.SynchronousDispatcher.invoke(SynchronousDispatcher.java:452)
        at org.jboss.resteasy.core.SynchronousDispatcher.lambda$invoke$4(SynchronousDispatcher.java:240)
        at org.jboss.resteasy.core.SynchronousDispatcher.lambda$preprocess$0(SynchronousDispatcher.java:154)
        at org.jboss.resteasy.core.interception.jaxrs.PreMatchContainerRequestContext.filter(PreMatchContainerRequestContext.java:321)
        at org.jboss.resteasy.core.SynchronousDispatcher.preprocess(SynchronousDispatcher.java:157)
        at org.jboss.resteasy.core.SynchronousDispatcher.invoke(SynchronousDispatcher.java:229)
        at io.quarkus.resteasy.runtime.standalone.RequestDispatcher.service(RequestDispatcher.java:82)
        at io.quarkus.resteasy.runtime.standalone.VertxRequestHandler.dispatch(VertxRequestHandler.java:147)
        at io.quarkus.resteasy.runtime.standalone.VertxRequestHandler$1.run(VertxRequestHandler.java:93)
        at io.quarkus.vertx.core.runtime.VertxCoreRecorder$14.runWith(VertxCoreRecorder.java:576)
        at org.jboss.threads.EnhancedQueueExecutor$Task.run(EnhancedQueueExecutor.java:2513)
        at org.jboss.threads.EnhancedQueueExecutor$ThreadBody.run(EnhancedQueueExecutor.java:1538)
        at org.jboss.threads.DelegatingRunnable.run(DelegatingRunnable.java:29)
        at org.jboss.threads.ThreadLocalResettingRunnable.run(ThreadLocalResettingRunnable.java:29)
        at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
        at java.base/java.lang.Thread.run(Thread.java:840)
Caused by: io.quarkus.arc.ArcUndeclaredThrowableException: Error invoking subclass method
        at se.sebank.portfoliomgt.tradingservice.trade.boundary.TradeResource_Subclass.validateTrades(Unknown Source)
        at se.sebank.portfoliomgt.tradingservice.trade.boundary.TradeResource_ClientProxy.validateTrades(Unknown Source)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
        at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.base/java.lang.reflect.Method.invoke(Method.java:568)
        at org.jboss.resteasy.core.MethodInjectorImpl.invoke(MethodInjectorImpl.java:154)
        at org.jboss.resteasy.core.MethodInjectorImpl.invoke(MethodInjectorImpl.java:118)
        at org.jboss.resteasy.core.ResourceMethodInvoker.internalInvokeOnTarget(ResourceMethodInvoker.java:560)
        at org.jboss.resteasy.core.ResourceMethodInvoker.invokeOnTargetAfterFilter(ResourceMethodInvoker.java:452)
        at org.jboss.resteasy.core.ResourceMethodInvoker.lambda$invokeOnTarget$2(ResourceMethodInvoker.java:413)
        at org.jboss.resteasy.core.interception.jaxrs.PreMatchContainerRequestContext.filter(PreMatchContainerRequestContext.java:321)
        at org.jboss.resteasy.core.ResourceMethodInvoker.invokeOnTarget(ResourceMethodInvoker.java:415)
        at org.jboss.resteasy.core.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:378)
        at org.jboss.resteasy.core.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:356)
        at org.jboss.resteasy.core.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:70)
        at org.jboss.resteasy.core.SynchronousDispatcher.invoke(SynchronousDispatcher.java:429)
        ... 15 more
Caused by: java.util.concurrent.ExecutionException: java.lang.NullPointerException: Cannot read the array length because "params" is null
        at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:396)
        at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2073)
        at se.sebank.portfoliomgt.tradingservice.trade.domain.TradeService.validateTrades(TradeService.java:125)
        at se.sebank.portfoliomgt.tradingservice.trade.domain.TradeService_ClientProxy.validateTrades(Unknown Source)
        at se.sebank.portfoliomgt.tradingservice.trade.boundary.TradeResource.validateTrades(TradeResource.java:74)

解决方案

1. 正确配置Mutiny实现并发执行

原代码串行的核心原因:transformToUniAndMerge默认并发度为1,且未给每个Uni指定工作线程池。修改如下:

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) {
    return Multi.createFrom().iterable(sendTrades)
            // 设置并发执行数量,根据系统资源调整
            .onItem().transformToUniAndMerge(5, tradeRequest ->
                    Uni.createFrom().item(() -> validateTrade(tradeRequest))
                            // 指定在自定义ManagedExecutor上执行,确保上下文传播
                            .runSubscriptionOn(managedExecutor)
            )
            .collect().asList()
            .await().atMost(Duration.ofMinutes(2));
}

2. 确保请求上下文与JWT传播

  • Quarkus会自动将请求上下文(包括JWT)传播到自身管理的线程池(如自定义ManagedExecutor或默认WorkerPool),禁止使用非Quarkus管理的线程池
  • 若validateTrade依赖RequestScoped Bean,需在方法上添加@ActivateRequestContext注解激活上下文:
@ActivateRequestContext
public Trade validateTrade(TradeRequest<TradeRequestDetails> tradeRequest) {
    // 原方法逻辑
}

3. Uni.combine的正确用法

如果选择用Uni.combine,需确保每个Uni都在工作线程池执行:

public List<Trade> validateTrades(final List<TradeRequest<TradeRequestDetails>> sendTrades) {
    List<Uni<Trade>> uniList = sendTrades.stream()
            .map(trade -> Uni.createFrom().item(() -> validateTrade(trade))
                    .runSubscriptionOn(managedExecutor))
            .toList();

    return Uni.combine().all().unis(uniList)
            .combinedWith(list -> list.stream().map(Trade.class::cast).toList())
            .await().atMost(Duration.ofMinutes(2));
}

内容的提问来源于stack exchange,提问作者Hassan Nazar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 18:55:56