如何阻塞等待异步StreamObserver调用执行完成?
解决gRPC异步StreamObserver阻塞等待响应完成的问题
最简单的实现方案:使用CountDownLatch
针对异步流场景下需要阻塞等待响应全部接收完成的需求,Java并发包中的CountDownLatch是最直接轻量的解决方案,比你提到的信号量更贴合单次等待流结束的场景。
修改后的代码示例:
// 初始化计数器,计数为1(等待响应流结束事件) CountDownLatch latch = new CountDownLatch(1); List<ServerReflectionResponse> responseList = new ArrayList<>(); StreamObserver<ServerReflectionResponse> responseStreamObserver = new StreamObserver<ServerReflectionResponse>() { @Override public void onNext(ServerReflectionResponse response) { responseList.add(response); } @Override public void onError(Throwable t) { // 错误场景必须触发计数器,避免线程永久阻塞 latch.countDown(); // 按需处理错误,例如打印日志或抛出异常 t.printStackTrace(); } @Override public void onCompleted() { // 响应全部接收完毕,触发计数器放行阻塞线程 latch.countDown(); } }; StreamObserver<ServerReflectionRequest> requestStreamObserver = reflectionStub.serverReflectionInfo(responseStreamObserver); requestStreamObserver.onNext(req); requestStreamObserver.onCompleted(); // 阻塞直到响应流结束或发生错误 try { latch.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 处理线程中断逻辑 } // 此时可安全读取列表数据 var sz = responseList.size(); // 后续可直接返回responseList或sz
额外说明
如果你的调用场景是单向请求-响应(而非双向流),可以直接使用gRPC的BlockingStub同步调用,代码会更简洁:
ServerReflectionResponse response = reflectionBlockingStub.serverReflectionInfo(req); // 直接获取响应结果,无需处理异步逻辑
但根据你使用的serverReflectionInfo方法定义(双向流),只能通过异步Stub配合阻塞工具实现同步等待。
内容的提问来源于stack exchange,提问作者DrPibb
相关产品推荐
相关产品推荐

