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

