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

Open Liberty Jakarta EE环境中gRPC StreamObserver线程缺失CDI上下文的问题咨询

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 19:09:34