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

Java实现双向gRPC重定向中间服务CANCELLED报错问题咨询

问题根因

你当前的实现核心错误在于没有遵循gRPC双向流的生命周期规则,具体问题如下:

  • 每收到一次上游(第一台服务器)的请求消息,就新建一次下游(第二台服务器)的gRPC调用,还在单条请求发送后立刻调用requestObserver.onCompleted()关闭了下游请求流,双向流在整个会话周期内需要复用同一个调用实例,不能单次消息就重建/关闭。
  • 下游响应流的onError、onCompleted事件没有同步转发给上游的响应观察者,上游流感知不到下游处理状态,超时后主动取消上下文,就会抛出你遇到的io.grpc.Context was cancelled报错。
  • 异步回调中没有绑定gRPC上下文,导致上下文丢失触发异常取消。
正确Java实现方案

Java完全可以实现该需求,不需要更换编程语言,修改后的代码如下:

package middle.server.pack;

import first.proto.pack.First;
import first.proto.pack.FirstProtoServiceGrpc;
import io.grpc.Context;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import second.proto.pack.Second;
import second.proto.pack.SecondProtoServiceGrpc;

import java.util.logging.LogManager;
import java.util.logging.Logger;

public class MiddleService extends FirstProtoServiceGrpc.FirstProtoServiceImplBase {
    // 下游连接全局复用,不要每次请求新建
    private final ManagedChannel channel = ManagedChannelBuilder.forTarget("localhost:8080").usePlaintext().build();
    private final SecondProtoServiceGrpc.SecondProtoServiceStub downstreamStub = SecondProtoServiceGrpc.newStub(channel);
    private final Logger logger = LogManager.getLogManager().getLogger(MiddleService.class.getName());

    @Override
    public StreamObserver<First.RequestToFirstServer> streamingCall(StreamObserver<First.ResponseForFirstServer> upstreamResponseObserver) {
        // 提前绑定当前gRPC上下文,避免异步回调丢失上下文
        Context currentContext = Context.current();

        // 整个流生命周期内只创建一次下游请求观察者
        StreamObserver<Second.RequestToSecondServer> downstreamRequestObserver = downstreamStub.streamingCall(
                // 包装下游响应观察者,绑定上下文
                currentContext.wrap(new StreamObserver<Second.ResponseFromSecondServer>() {
                    @Override
                    public void onNext(Second.ResponseFromSecondServer value) {
                        doProcessOnResponse(value);
                        First.ResponseForFirstServer response = mapToFirstResponse(value);
                        upstreamResponseObserver.onNext(response);
                    }

                    @Override
                    public void onError(Throwable t) {
                        logger.info("下游调用出错: " + t.getMessage());
                        // 错误同步转发给上游
                        upstreamResponseObserver.onError(t);
                    }

                    @Override
                    public void onCompleted() {
                        logger.info("下游处理完成");
                        // 完成事件同步转发给上游
                        upstreamResponseObserver.onCompleted();
                    }
                })
        );

        // 返回上游请求观察者,所有事件直接转发给下游
        return new StreamObserver<First.RequestToFirstServer>() {
            @Override
            public void onNext(First.RequestToFirstServer value) {
                Second.RequestToSecondServer downstreamRequest = mapToSecondRequest(value);
                downstreamRequestObserver.onNext(downstreamRequest);
            }

            @Override
            public void onError(Throwable t) {
                logger.info("上游请求出错: " + t.getMessage());
                // 上游错误同步通知下游关闭流
                downstreamRequestObserver.onError(t);
            }

            @Override
            public void onCompleted() {
                logger.info("上游请求发送完成");
                // 上游请求完成,通知下游不再发送请求
                downstreamRequestObserver.onCompleted();
            }
        };
    }

    // 以下方法需自行实现业务逻辑
    private void doProcessOnResponse(Second.ResponseFromSecondServer value) {
        // 自定义响应处理逻辑
    }

    private First.ResponseForFirstServer mapToFirstResponse(Second.ResponseFromSecondServer value) {
        return First.ResponseForFirstServer.newBuilder()
                .setSomeprocessedinformation(value.getComputedInformation())
                .build();
    }

    private Second.RequestToSecondServer mapToSecondRequest(First.RequestToFirstServer value) {
        Second.RequestToSecondServer.Builder builder = Second.RequestToSecondServer.newBuilder();
        if (value.hasX()) {
            builder.setProcessedX(value.getX());
        } else if (value.hasY()) {
            builder.setProcesdedY(value.getY());
        }
        return builder.build();
    }

    // 服务停止时记得关闭下游连接,避免资源泄漏
    public void shutdown() throws InterruptedException {
        channel.shutdown().awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS);
    }
}
注意事项
  • 如果你连接的下游服务没有启用TLS,要记得加.usePlaintext()配置,否则会连接失败。
  • 生产环境建议配置下游连接的超时、重试策略,避免单点故障影响服务可用性。

内容的提问来源于stack exchange,提问作者Нуркасым Абилов

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 06:39:03