【问题标题】:Graceful shutdown of gRPC downstreamgRPC 下游的优雅关闭
【发布时间】:2018-08-14 07:39:15
【问题描述】:

使用以下 proto 缓冲区代码:

syntax = "proto3";

package pb;

message SimpleRequest {
    int64 number = 1;
}

message SimpleResponse {
    int64 doubled = 1;
}

// All the calls in this serivce preform the action of doubling a number.
// The streams will continuously send the next double, eg. 1, 2, 4, 8, 16.
service Test {
    // This RPC streams from the server only.
    rpc Downstream(SimpleRequest) returns (stream SimpleResponse);
}

我能够成功打开一个流,并不断从服务器获取下一个双倍数字。

我的运行代码如下:

ctxDownstream, cancel := context.WithCancel(ctx)
downstream, err := testClient.Downstream(ctxDownstream, &pb.SimpleRequest{Number: 1})
for {
    responseDownstream, err := downstream.Recv()
    if err != io.EOF {
        println(fmt.Sprintf("downstream response: %d, error: %v", responseDownstream.Doubled, err))

        if responseDownstream.Doubled >= 32 {
            break
        }
    }
}
cancel() // !!This is not a graceful shutdown
println(fmt.Sprintf("%v", downstream.Trailer()))

我遇到的问题是使用上下文取消意味着我的下游.Trailer() 响应为空。有没有办法从客户端优雅地关闭这个连接并接收下游.Trailer()。

注意:如果我从服务器端关闭下游连接,我的预告片就会被填充。但是我无法指示我的服务器端关闭这个特定的流。所以必须有一种方法可以优雅地关闭流客户端。

谢谢。

根据要求提供一些服务器代码:

func (b *binding) Downstream(req *pb.SimpleRequest, stream pb.Test_DownstreamServer) error {
    request := req

    r := make(chan *pb.SimpleResponse)
    e := make(chan error)
    ticker := time.NewTicker(200 * time.Millisecond)
    defer func() { ticker.Stop(); close(r); close(e) }()

    go func() {
        defer func() { recover() }()
        for {
            select {
            case <-ticker.C:
                response, err := b.Endpoint(stream.Context(), request)
                if err != nil {
                    e <- err
                }
                r <- response
            }
        }
    }()

    for {
        select {
        case err := <-e:
            return err
        case response := <-r:
            if err := stream.Send(response); err != nil {
                return err
            }
            request.Number = response.Doubled
        case <-stream.Context().Done():
            return nil
        }
    }
}

您仍然需要在预告片中填充一些信息。我使用 grpc.StreamServerInterceptor 来执行此操作。

【问题讨论】:

  • 您能否也发布您的服务器实现代码?
  • 我提供了一些服务器代码
  • 你能提供一个最小的 git repo 来调试吗?它可能会更快地为您提供解决方案
  • 出于兴趣,我的回答对您不起作用有什么原因吗?

标签: go grpc


【解决方案1】:

根据 grpc go 文档

Trailer 从服务器返回预告片元数据(如果有)。 它只能在 stream.CloseAndRecv 返回后调用,或者 stream.Recv 返回了一个非零错误(包括 io.EOF)

所以如果你想在客户端阅读预告片,试试这样的

ctxDownstream, cancel := context.WithCancel(ctx)
defer cancel()
for {
  ...
  // on error or EOF
  break;
}
println(fmt.Sprintf("%v", downstream.Trailer()))

出现错误时退出无限循环并打印预告片。 cancel 将在函数结束时被调用,因为它被延迟了。

【讨论】:

  • 当我这样做时,预告片总是空的
  • @Jamie 你能分享一下你是如何在拦截器中添加预告片的代码吗?
【解决方案2】:

我找不到能清楚解释它的参考资料,但这似乎是不可能的。

在线上,grpc-status 后跟在调用正常完成时(即服务器退出调用)时的预告元数据。 当客户端取消调用时,这些都不会发送。

似乎 gRPC 将调用取消视为 rpc 的快速中止,与丢弃的套接字没有太大区别。

通过请求流作品添加“取消消息”;服务器可以拾取它并从其末端取消流,并且仍然会发送预告片:

message SimpleRequest {
    oneof RequestType {
        int64 number = 1;
        bool cancel = 2;
    }
}
....
rpc Downstream(stream SimpleRequest) returns (stream SimpleResponse);

虽然这确实给代码增加了一点复杂性。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-02-05
    • 1970-01-01
    • 1970-01-01
    • 2018-02-05
    • 1970-01-01
    • 2019-09-26
    • 2019-05-16
    • 1970-01-01
    相关资源
    最近更新 更多