gRPC-Java RPC 取消(Cancellation)机制详解:Context 与 ServerCallStreamObserver 双 API 实战
gRPC-Java RPC 取消Cancellation机制详解Context 与 ServerCallStreamObserver 双 API 实战【免费下载链接】grpc-javaThe Java gRPC implementation. HTTP/2 based RPC项目地址: https://gitcode.com/GitHub_Trending/gr/grpc-javagRPC 客户端在不再关心某个 RPC 调用结果时可以主动取消该调用向服务端传达不再关注的信号。本文以 grpc-java 仓库中的 Cancellation 示例 为骨架完整讲解取消的定义与触发原因、服务端感知取消的两套 APIio.grpc.Context与ServerCallStreamObserver、以及客户端在不同 Stub 形态下发起取消的四种方式。读完本文你将掌握在阻塞、Future、异步 Stub 下如何正确取消 RPC并能根据服务端是同步还是异步实现选择合适的中止通知机制。什么是 RPC 取消Cancellation按照示例文档 Cancellation README 与 CancellationServer.java 的定义任何对正在进行的 RPC 的中止abort都被视为该 RPC 的取消。常见的取消诱因有以下三类客户端显式取消客户端代码主动调用取消 API例如context.cancel()、future.cancel(true)、StreamObserver.onError()等截止时间Deadline到期RPC 设置了 deadline超时后由客户端/传输层自动触发取消I/O 故障网络中断、连接断开等传输层错误也会导致 RPC 中止。一个关键事实是服务端永远不会被告知取消的具体原因。也就是说服务端只知道这个 RPC 被取消了而无法得知是客户端主动放弃、deadline 超时还是网络故障。这一点在示例文档和 ServerCallStreamObserver 的 Javadoc 中均有明确表述取消可能由超时、客户端显式取消、网络错误等原因造成。服务端感知取消的两套 API服务端有两种 API 可以获知 RPC 被取消二者在回调线程模型上存在本质差异需要根据服务端实现是同步还是异步来权衡。API回调线程线程安全要求适用场景io.grpc.ContextCancellationListener回调在另一个线程执行监听器必须自行保证线程安全同步/阻塞式服务实现如unaryEchoServerCallStreamObserversetOnCancelHandler与其它StreamObserver回调串行化执行应用通常无需额外线程安全处理真正异步的服务实现如基于调度器的流式bidirectionalStreamingEcho两套 API 都提供线程安全的isCancelled()轮询方法可随时检查当前 RPC 是否已被取消。基于 ServerCallStreamObserver 的取消通知ServerCallStreamObserver位于 stub 模块其核心 API 为boolean isCancelled()当调用被取消且服务端应中止处理以节省资源时返回true可以安全地被多个线程并发调用源码第 44 行void setOnCancelHandler(Runnable onCancelHandler)注册取消回调源码第 67 行。该Runnable的执行与入站StreamObserver回调严格串行化——即不会与onNext/onCompleted/onError并发执行正因如此若其它回调运行时间较长取消回调会被延迟此时可改用不会被延迟的isCancelled()轮询。注意事项源码 Javadoc 明确给出setOnCancelHandler只能在服务方法首次被调用期间、服务端返回其StreamObserver之前注册设置onCancelHandler会抑制onNext()中抛出的取消异常若调用方已通过轮询处理取消、无需实质观察取消事件使用 no-op 的onCancelHandler也可以只用来抑制该异常。基于 Context 的取消通知io.grpc.Context的CancellationListener回调运行在另一个线程因此监听器代码必须线程安全。示例服务端在unaryEcho中的典型用法是结合一个可取消的FutureTask通过监听器在 RPC 被取消时中断/取消正在进行的阻塞操作见下文完整代码。示例运行环境Echo 服务与构建方式Cancellation 示例基于仓库内置的 Echo 服务其 proto 定义位于 echo.proto包含四种 RPC 形态UnaryEcho一元ServerStreamingEcho服务端流ClientStreamingEcho客户端流BidirectionalStreamingEcho双向流Cancellation 示例只使用其中的UnaryEcho与BidirectionalStreamingEcho恰好分别演示同步与异步两种服务端实现下的取消处理。构建与运行方式详见 examples/README.md# 在 grpc-java/examples 目录下非发布版本需先按 COMPILING.md 本地安装 SNAPSHOT $ ./gradlew installDist构建成功后启动脚本会生成在build/install/examples/bin/下。Cancellation 示例的两个启动脚本由 examples/build.gradle 注册$ ./build/install/examples/bin/cancellation-server # 终端一先启动服务端默认监听 50051 $ ./build/install/examples/bin/cancellation-client # 终端二再启动客户端客户端支持通过命令行参数指定目标地址默认localhost:50051传--help可查看用法见 CancellationClient.java。服务端实现两种取消感知方式的完整代码服务端 CancellationServer.java 内部类SlowEcho继承EchoGrpc.EchoImplBase分别用两种方式处理取消。双向流用 ServerCallStreamObserver 感知取消对于真正异步的实现使用ServerCallStreamObserver接收取消通知效果最好。服务方法中可安全地把传入的responseObserver强转为ServerCallStreamObserverServerCallStreamObserverEchoResponse responseCallObserver (ServerCallStreamObserverEchoResponse) responseObserver;服务端每收到一个请求就用单线程调度器按固定速率每 200ms持续回显同时通过setOnCancelHandler注册取消回调在取消发生时停止所有定时回显任务responseCallObserver.setOnCancelHandler(requestObserver::onCancel); // ... public void onCancel() { // 若此时 onCompleted() 尚未被调用则本方法与 onError 都会被调用 // 若 onCompleted() 已被调用则只有本方法被调用。 System.out.println(Bidi RPC cancelled); stopEchos(echos); }这里还有一个重要细节源码注释明确说明onCancel()即使在服务端已经完成或失败该 RPC 之后仍可能被调用因为回调存在竞态、响应仍需发送给客户端。如果需要精确区分RPC 正常完成与被取消应改用setOnCloseHandler()来获得服务端尽力而为的完成通知。一元 RPC用 Context 感知取消setOnCancelHandler对unaryEcho并不适用因为该方法只有在产出结果后才会返回而ServerCallStreamObserver保证Runnable不会与其它 RPC 回调包括本方法并发执行取消通知注定来得太晚。因此示例采用Context方案Context currentContext Context.current(); for (int i 0; i 10; i) { if (currentContext.isCancelled()) { System.out.println(Unary RPC cancelled); responseObserver.onError( Status.CANCELLED.withDescription(RPC cancelled).asRuntimeException()); return; } FutureTaskVoid task new FutureTask(() - { Thread.sleep(100); // Do some work return null; }); // 通过 Context 监听器在 RPC 被取消时从另一个线程中断正在进行的操作 Context.CancellationListener listener (Context context) - task.cancel(true); Context.current().addListener(listener, MoreExecutors.directExecutor()); task.run(); // 一个可取消的操作 Context.current().removeListener(listener); }关键语义差异源码注释明确指出ServerCallStreamObserver.isCancelled()只有 RPC被取消时才返回trueContext.isCancelled()类似但 RPC正常完成时也会返回true对示例而言用哪个 API 轮询都无所谓但理解这一差异对生产代码很关键。此外gRPC Stub 会观察io.grpc.Context的取消状态因此进行嵌套 RPC 时取消会自动传播。如果需要禁用这种自动传播可以换用不同的 Context或使用Context.ROOT.call(...)/Context.fork().call(...)例如Context.ROOT.call(() - futureStub.unaryEcho(request)); // 禁用传播 context.fork().call(() - futureStub.unaryEcho(request)); // 分叉 Context客户端取消的四种方式客户端 CancellationClient.java 的demonstrateCancellation()依次演示了四种取消手法覆盖了阻塞、Future、异步三种 Stub 形态。方式一阻塞 Stub CancellableContextio.grpc.Context可以配合任意Stub 取消 RPC也是唯一能取消阻塞 Stub RPC的方式。它本质上是一种通用的、可替代线程中断thread interruption的机制也可在 gRPC 之外使用用于应用内部的协调。// CancellableContext 必须在其生命周期结束时被 cancel 或 close否则可能造成内存泄漏 try (CancellableContext context Context.current().withCancellation()) { new Thread(() - { try { Thread.sleep(500); // 做一些工作 } catch (InterruptedException ex) { Thread.currentThread().interrupt(); } // 取消原因永远不会发送给服务端但会作为 RPC 失败原因回显给客户端 context.cancel(new RuntimeException(Oops. Messed that up, let me try again)); }).start(); // context.run() 把该 Context 附着到当前线程供 gRPC 观察返回前会自动恢复之前的 Context context.run(() - echoBlocking(RAAWRR haha lol hehe AWWRR GRRR)); }要点CancellableContext必须用 try-with-resources或手动 close/cancel释放否则会泄漏内存context.run()负责挂载与恢复 Context。方式二Future Stub future.cancel(true)用cancel(true)取消ListenableFuture会连带取消底层 RPCListenableFutureEchoResponse future echoFuture(Future clie*cough*nt was here!); Thread.sleep(500); // 做一些工作 future.cancel(true); // 突然不想要这个回显了 Thread.sleep(100); // 让日志更明显取消是异步的echoFuture内部使用EchoGrpc.newFutureStub(channel).unaryEcho(request)并通过Futures.addCallback注册成功/失败回调失败时用Status.fromThrowable(t)打印状态。方式三异步 Stub StreamObserver.onError()异步的onError()会触发取消但不能与 StreamObserver 上的其它调用并发执行。若需要线程安全应改用方式一的CancellableContextStreamObserverEchoRequest reqObserver echoAsync(... async client... is the... best...); Thread.sleep(500); // 做一些工作 // 只要还没调用 onCompleted()就可以用 onError() 取消 reqObserver.onError(new RuntimeException(That was weak...));方式四异步 Stub ClientCallStreamObserver.cancel()ClientCallStreamObserver也提供cancel(String message, Throwable cause)方法。它同样不允许与其它调用并发reqCallObserver echoAsync(Async client or bust!); reqCallObserver.onCompleted(); Thread.sleep(250); // 做一些工作 // 既然已经调用了 onCompleted()就不能再用 onError()此时用 cancel() 是安全的 reqCallObserver.cancel(Thats enough. Im bored, null);关于ClientCallStreamObserver的获取方式源码注释给出了关键指引CancellationClient.java客户端流与双向流Stub 方法返回的StreamObserver可直接强转为ClientCallStreamObserver一元与服务器流Stub 方法不返回StreamObserver需改用ClientResponseObserver来获得ClientCallStreamObserver由于ClientCallStreamObserver.cancel()不是线程安全的在 RPC Stub 方法如unaryEcho()返回之前不能从其它线程调用它。取消是异步的时序与竞态要点从示例代码的多处Thread.sleep(...)与注释可以看出取消Cancel是异步操作客户端调用取消 API 后服务端感知到取消并做出反应存在时间差。因此示例中在每次取消后都预留了短暂延时让日志顺序更可读。同时要注意两个竞态来源服务端回调的串行化保证setOnCancelHandler的回调虽然与其它回调串行但如果其它回调执行时间过长取消通知会被延迟此时用isCancelled()轮询可绕过延迟onCancel()可能在 RPC 完成后仍被调用需要精确区分完成与取消时使用setOnCloseHandler()而非仅依赖onCancel()。小结gRPC-Java 的取消机制设计围绕一个核心事实展开任何 RPC 中止都是取消且服务端不会获知原因。服务端要感知取消同步阻塞场景选io.grpc.Context注意监听器线程安全与Context.isCancelled()在正常完成时也返回 true 的语义异步流式场景选ServerCallStreamObserverisCancelled()轮询 setOnCancelHandler回调客户端则按 Stub 形态在CancellableContext、future.cancel(true)、onError()、ClientCallStreamObserver.cancel()四者之间取舍。结合 CancellationClient.java 与 CancellationServer.java 两个完整示例即可在真实代码中正确实现取消的发起、传播与资源回收。【免费下载链接】grpc-javaThe Java gRPC implementation. HTTP/2 based RPC项目地址: https://gitcode.com/GitHub_Trending/gr/grpc-java创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考