Open Liberty Jakarta EE环境中gRPC StreamObserver线程缺失CDI上下文的问题咨询
我目前在运行Open Liberty 23.0.0.9服务器,基于Jakarta 10 + CDI 3.0环境,现在通过gRPC通道和外部服务器建立了推拉模式的数据流。我给gRPC通道指定了ManagedServiceExecutor,按我的理解,这样线程应该由容器管理,并且带有活跃的CDI上下文,但实际运行下来却不是这么回事,或者说我的理解存在偏差。
我是这样创建gRPC通道的:
ManagedChannelBuilder.forAddress(grpcHost, port) .enableRetry() .keepAliveWithoutCalls(true) .executor(executorService) .build(); PubSubGrpc.newStub(channel).withCallCredentials(credentials).subscribe(responseObserver);
对应的ResponseObserver实现如下:
public StreamObserver<FetchResponse> getDefaultResponseStreamObserver() { return new StreamObserver<FetchResponse>() { @Override public void onNext(FetchResponse fetchResponse) { for (ConsumerEvent ce : fetchResponse.getEventsList()) { try { injectedBean.methodThatRequiresCDI(ce); // 这个方法无法正常工作 } catch (Exception e) { logger.info(e.toString()); } } if (fetchResponse.getPendingNumRequested() == 0 && subscriptionReference.isActive()) { subscriptionReference.getRequestObserver().onNext(FetchRequest.newBuilder().setNumRequested(100).build()); } } @Override public void onError(Throwable t) { subscriptionReference = subscriptionReference.closedSubscription(); logger.info("big bad error :("); } @Override public void onCompleted() { subscriptionReference = subscriptionReference.closedSubscription(); logger.info("Call completed by server. Closing AsyncSubscriptionFactory. Goodbye."); } };
问题出在injectedBean.methodThatRequiresCDI(ce)这行:这个方法在下游依赖库中调用了CDI.current(),但会抛出IllegalStateException。为了简化问题,我直接在onNext()方法里调用CDI.current(),也会出现完全相同的异常和堆栈信息——这说明当gRPC流触发onNext()回调时,当前线程根本没有活跃的CDI上下文。
异常堆栈信息如下:
java.lang.IllegalStateException: Could not find deployment at com.ibm.ws.cdi.impl.AbstractCDIRuntime.getCDI(AbstractCDIRuntime.java:163) at jakarta.enterprise.inject.spi.CDI.getCDIProvider(CDI.java:78) at jakarta.enterprise.inject.spi.CDI.current(CDI.java:65) at ----.subscribe.control.ExampleSubscriber$1.onNext(ExampleSubscriber.java:107) at ----.subscribe.control.ExampleSubscriber$1.onNext(ExampleSubscriber.java:100) at io.grpc.stub.ClientCalls$StreamObserverToCallListenerAdapter.onMessage(ClientCalls.java:466) at io.grpc.ForwardingClientCallListener.onMessage(ForwardingClientCallListener.java:33) at io.openliberty.grpc.internal.client.monitor.GrpcMonitoringClientCallListener.onMessage(GrpcMonitoringClientCallListener.java:59) at io.grpc.ForwardingClientCallListener.onMessage(ForwardingClientCallListener.java:33) at io.grpc.ForwardingClientCallListener.onMessage(ForwardingClientCallListener.java:33) at io.grpc.ForwardingClientCallListener.onMessage(ForwardingClientCallListener.java:33) at io.grpc.internal.DelayedClientCall$DelayedListener.onMessage(DelayedClientCall.java:447) at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInternal(ClientCallImpl.java:661) at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInContext(ClientCallImpl.java:646) at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133) at com.ibm.ws.threading.internal.PolicyTaskFutureImpl.run(PolicyTaskFutureImpl.java:762) at com.ibm.ws.threading.internal.PolicyExecutorImpl.runTask(PolicyExecutorImpl.java:1172) at com.ibm.ws.threading.internal.PolicyExecutorImpl$GlobalPoolTask.run(PolicyExecutorImpl.java:198) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:834)
需要说明的是,这个methodThatRequiresCDI方法在其他上下文里运行完全正常——比如作为JAX-RS请求的一部分,或者由MDB调用时,只要CDI.current()能正常获取上下文、BeanManager可访问,方法调用就没有任何问题。
目前我找到了两个可行的临时解决方案,但都不是最优的:
- 方案一:把持有
methodThatRequiresCDI方法的Bean改成EJB@Singleton。这样调用时会在容器管理的线程中执行,CDI.current()就能正常工作。但现在这个Bean是@ApplicationScoped的,我希望保持这个作用域;而且我需要调用很多不同的Bean,不想把它们都改成EJB Singleton。 - 方案二:把消息负载包装成事件触发,用带有
@Observes注解的方法来间接调用目标方法。这种方式也能正常工作,CDI.current()在观察者方法里能正常执行,下游方法也没问题。但缺点是多了一层没必要的间接调用,感觉像是为了让代码跑起来的hack;而且如果有很多不同类型的订阅,重构起来会很麻烦,每种订阅都要单独触发事件并添加对应的@Observes方法。
我更想搞清楚的是:为什么调用onNext()的线程没有CDI上下文?有没有办法直接解决这个问题?我原本以为只要使用容器创建的线程、把任务提交给ManagedExecutorService,就会自动带上CDI上下文。我也试过用ManagedThreadFactory、在onNext()里把任务提交给ManagedThreadManager、添加@ActivateRequestContext注解等方法,但都没有效果。
备注:内容来源于stack exchange,提问作者daniel strandberg

