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
相关产品推荐
相关产品推荐

