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

gRPC Fire-and-Forget模式实现遇CANCELLED异常求助

解决gRPC Fire-and-Forget模式下的CANCELLED异常问题

异常原因分析

你遇到的io.grpc.StatusRuntimeException: CANCELLED: io.grpc.Context was cancelled without error异常,本质是因为异步gRPC调用默认绑定当前方法的Context。当sendCommit方法执行完毕返回后,当前线程的Context会被取消,导致正在进行的异步RPC请求被强制终止。

实现Fire-and-Forget的正确方案

要实现真正的即发即弃模式,需要让RPC请求的生命周期脱离当前方法的Context,确保请求能在后台独立完成。以下是两种可行的实现方式:

方式一:使用独立的Context

通过Context.fork()创建一个独立的Context,在这个Context中执行RPC调用,避免受当前方法Context的影响:

private Boolean sendCommit(List<Transaction> allTransactions) {
    for (int cur_server_number = 1; cur_server_number <= totalServers; cur_server_number++) {
        if (cur_server_number == SERVER_NUMBER) {
            continue;
        }
        try{
            ManagedChannel channel = serverChannels[cur_server_number];
            DistributedBankGrpc.DistributedBankStub stub = DistributedBankGrpc.newStub(channel);
            BallotData ballot = BallotData.newBuilder()
                    .setBallotNum(this.ballotNumber)
                    .setServerNum(SERVER_NUMBER)
                    .build();
            CommitRequest commitRequest = CommitRequest.newBuilder()
                    .setBallot(ballot)
                    .addAllTransactions(allTransactions)
                    .build();

            final int serverNumber = cur_server_number;
            // 创建独立的Context并执行RPC
            Context.fork().run(() -> {
                stub.commit(commitRequest, new StreamObserver<Empty>() {
                    @Override
                    public void onNext(Empty emptyResponse) {
                        System.out.println("Commited on server " + serverNumber);
                    }

                    @Override
                    public void onError(Throwable t) {
                        System.err.println("Commit request to Server " + serverNumber + " failed: " + t.getMessage());
                        t.printStackTrace();
                    }

                    @Override
                    public void onCompleted() {
                        System.out.println("Commit request to Server " + serverNumber + " completed.");
                    }
                });
            });
        } catch (Exception e) {
            System.out.println("An exception occurred: " + e.getMessage());
        }
    }

    return true;
}

方式二:使用线程池异步提交任务

将RPC调用提交到独立的线程池执行,让请求在后台线程中运行,彻底脱离当前方法的上下文:

// 提前初始化一个线程池(可根据业务需求调整参数)
private static final ExecutorService RPC_EXECUTOR = Executors.newFixedThreadPool(10);

private Boolean sendCommit(List<Transaction> allTransactions) {
    for (int cur_server_number = 1; cur_server_number <= totalServers; cur_server_number++) {
        if (cur_server_number == SERVER_NUMBER) {
            continue;
        }
        try{
            ManagedChannel channel = serverChannels[cur_server_number];
            DistributedBankGrpc.DistributedBankStub stub = DistributedBankGrpc.newStub(channel);
            BallotData ballot = BallotData.newBuilder()
                    .setBallotNum(this.ballotNumber)
                    .setServerNum(SERVER_NUMBER)
                    .build();
            CommitRequest commitRequest = CommitRequest.newBuilder()
                    .setBallot(ballot)
                    .addAllTransactions(allTransactions)
                    .build();

            final int serverNumber = cur_server_number;
            // 提交到线程池执行RPC
            RPC_EXECUTOR.submit(() -> {
                stub.commit(commitRequest, new StreamObserver<Empty>() {
                    @Override
                    public void onNext(Empty emptyResponse) {
                        System.out.println("Commited on server " + serverNumber);
                    }

                    @Override
                    public void onError(Throwable t) {
                        System.err.println("Commit request to Server " + serverNumber + " failed: " + t.getMessage());
                        t.printStackTrace();
                    }

                    @Override
                    public void onCompleted() {
                        System.out.println("Commit request to Server " + serverNumber + " completed.");
                    }
                });
            });
        } catch (Exception e) {
            System.out.println("An exception occurred: " + e.getMessage());
        }
    }

    return true;
}

关键注意事项

  • ManagedChannel复用:确保serverChannels中的Channel是全局复用的,不要每次调用都创建新的Channel,否则会导致资源泄漏。
  • 错误处理:即使是即发即弃模式,也要保留onError的日志记录,便于后续排查问题。
  • 线程池资源管理:如果使用线程池方式,需要在应用关闭时调用RPC_EXECUTOR.shutdown()释放线程资源。

内容的提问来源于stack exchange,提问作者The Beast

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 13:36:04