【问题标题】:How to design publish-subscribe pattern properly in grpc?如何在 grpc 中正确设计发布订阅模式?
【发布时间】:2019-07-23 08:17:12
【问题描述】:

我正在尝试使用 grpc 实现 pub sub 模式,但我对如何正确执行有点困惑。

我的原型:rpc call (google.protobuf.Empty) returns (stream Data);

客户:

asynStub.call(Empty.getDefaultInstance(), new StreamObserver<Data>() {
         @Override
         public void onNext(Data value) {
           // process a data

         @Override
         public void onError(Throwable t) {

         }

         @Override
         public void onCompleted() {

         }
       });

   } catch (StatusRuntimeException e) {
     LOG.warn("RPC failed: {}", e.getStatus());
   }

   Thread.currentThread().join();

服务器服务:

public class Sender extends DataServiceGrpc.DataServiceImplBase implements Runnable {
  private final BlockingQueue<Data> queue;
  private final static HashSet<StreamObserver<Data>> observers = new LinkedHashSet<>();

  public Sender(BlockingQueue<Data> queue) {
    this.queue = queue;
  }

  @Override
  public void data(Empty request, StreamObserver<Data> responseObserver) {
    observers.add(responseObserver);
  }

  @Override
  public void run() {
    while (!Thread.currentThread().isInterrupted()) {
      try {
        // waiting for first element
        Data data = queue.take();
        // send head element
        observers.forEach(o -> o.onNext(data));

      } catch (InterruptedException e) {
        LOG.error("error: ", e);
        Thread.currentThread().interrupt();
      }
    }
  }
}

如何正确地从全局观察者中移除客户端?连接断开时如何接收某种信号?
如何管理客户端-服务器重新连接?连接断开时如何强制客户端重新连接?

提前致谢!

【问题讨论】:

    标签: java publish-subscribe grpc


    【解决方案1】:

    在您的服务实施中:

      @Override
      public void data(Empty request, StreamObserver<Data> responseObserver) {
        observers.add(responseObserver);
      }
    

    需要获取当前请求的Context,和listen for cancellation。对于单请求、多响应调用(也称为服务器流),gRPC 生成的代码被简化为直接传递请求。这意味着您无法直接访问底层的ServerCall.Listener,而这正是您通常监听客户端断开连接和取消的方式。

    相反,每个 gRPC 调用都有一个与之关联的Context,它携带取消和其他请求范围的信号。对于您的情况,您只需要通过添加自己的侦听器来侦听取消,然后从链接的哈希集中安全地删除响应观察器。


    关于重连:gRPC 客户端会在连接中断时自动重连,但通常不会重试 RPC,除非这样做是安全的。在服务器流式 RPC 的情况下,这样做通常不安全,因此您需要直接在客户端重试 RPC。

    【讨论】:

    • 卡尔,谢谢!在 grpc 之上构建服务器流(发布/订阅)是否是一种好习惯?
    • 是的。 Google 的 Cloud Pubsub 位于 gRPC 之上,其客户端源代码在 GitHub 上是公开的。你可以看看它以获得灵感。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多