【问题标题】:grpc server stops receiving messages after sending many messages simultaneouslygrpc 服务器同时发送多条消息后停止接收消息
【发布时间】:2019-07-29 12:22:07
【问题描述】:

我正在实现一个简单的 grpc 服务,其中任务的摘要将被发送到 grpc 服务器。如果我发送的消息数量较少,一切正常,但是当我开始发送 5000 条消息时,服务器停止并在客户端收到超出期限的消息。我也尝试重新连接,但发现错误消息为。

rpc error: code = Unavailable desc = all SubConns are in TransientFailure, latest connection error: timed out waiting for server handshake

服务器没有显示错误并且是活动的。

我也尝试设置 GRPC_GO_REQUIRE_HANDSHAKE=off,但错误仍然存​​在。我还实现了批量发送摘要,但重复相同的场景。

在 grpc 中发送的消息数量有限制吗?

这是我的服务原型


// The Result service definition.
service Result {
  rpc ConntectMaster(ConnectionRequest) returns (stream ExecutionCommand) {}
  rpc postSummary(Summary) returns(ExecutionCommand) {}
}

message Summary{

  int32 successCount = 1;
  int32 failedCount = 2;
  int32 startTime = 3;
  repeated TaskResult results = 4;
  bool isLast = 5;
  string id = 6;
}

服务器中的postSummary实现

// PostSummary posts the summary to the master
func (server *Server) PostSummary(ctx context.Context, in *pb.Summary) (*pb.ExecutionCommand, error) {

    for i := 0; i < len(in.Results); i++ {

        res := in.Results[i]
        log.Printf("%s --> %d Res :: %s, len : %d", in.Id, i, res.Id, len(in.Results))

    }
    return &pb.ExecutionCommand{Type: stopExec}, nil
}
func postSummaryInBatch(executor *Executor, index int) {
    summary := pb.Summary{
        SuccessCount: int32(executor.summary.successCount),
        FailedCount:  int32(executor.summary.failedCount),
        Results:      []*pb.TaskResult{},
        IsLast:       false,
    }

    if index >= len(executor.summary.TaskResults) {
        summary.IsLast = true
        return
    }

    var to int
    batch := 500
    if (index + batch) <= len(executor.summary.TaskResults) {
        to = index + batch
    } else {
        to = len(executor.summary.TaskResults)
    }
    for i := index; i < to; i++ {
        result := executor.summary.TaskResults[i]
        taskResult := pb.TaskResult{
            Id:   result.id,
            Msg:  result.msg,
            Time: result.time,
        }
        // log.Printf("adding res : %s ", taskResult.Id)

        if result.err != nil {
            taskResult.IsError = true
        }
        summary.Results = append(summary.Results, &taskResult)
    }
    summary.Id = fmt.Sprintf("%d-%d", index, to)
    log.Printf("sent from  %d to %d ", index, to)
    postSummary(executor, &summary, 0)
    postSummaryInBatch(executor, to)
}

func postSummary(executor *Executor, summary *pb.Summary, retryCount int) {
    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()

    cmd, err := client.PostSummary(ctx, summary)
    if err != nil {
        if retryCount < 3 {
            reconnect(executor)
            postSummary(executor, summary, retryCount+1)
        }
        log.Printf(err.Error())
        // log.Fatal("cannot send summary report")
    } else {
        processServerCommand(executor, cmd)
    }
}

【问题讨论】:

    标签: go grpc


    【解决方案1】:

    grpc 默认 maxReceiveMessageSize 为 4MB,您的 grpc 客户端可能超过了该限制。

    grpc 在传输层中使用 h2,它只打开一个 tcp 连接并在其上多路复用“请求”,与 h1 相比减少了显着的开销,我不会太担心批处理,只会对 grpc 服务器进行单独调用。

    【讨论】:

    • 我将服务器和客户端的 maxReceiveMessageSize 以及发送消息大小更改为 20MB,但问题仍然存在。
    • 我猜您对消息长度的看法是正确的,因为现在当我不发送摘要结果列表时,消息已成功发送。如果我需要发送结果列表,我应该怎么做。我也尝试过发送流,但没有成功。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-06
    • 2011-11-23
    相关资源
    最近更新 更多