【问题标题】:Get streamed file size from grpc server interceptor从 grpc 服务器拦截器获取流文件大小
【发布时间】:2021-07-28 06:27:48
【问题描述】:

我在 proto-file 中有一个像这样定义的 Go 服务器双向流方法:

syntax = "proto3";
option go_package="pdfcompose/;pdfcompose";
package pdfcompose;

service PdfCompose {
  rpc Send (stream FileForm) returns (stream PdfFile) {}
}

message FileForm {
  bytes Upfile1 = 1;
  bytes Upfile2 = 2;
  bytes Upfile3 = 3;
}

message PdfFile {
  bytes File = 1;
}

而我的日志拦截器有如下接口:

func logInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler)  error {
    fmt.Println("Log Interceptor")
    err := handler(srv, ss)
    if err != nil {
        return err
    }
    return nil
}

我使用https://github.com/grpc-ecosystem/go-grpc-middleware 作为拦截器引擎。 我需要实现流文件大小的日志记录(出于教育目的),并试图找出我可以从哪里获取有关 FileForm 及其内容的任何数据。

我的第一个猜测是查看 grpc.ServerStream 参数 (ss) 以查找有关它的信息,它看起来包含大量数据,例如最大和最小 MessageSize,但注意实际内容长度。

如何使用这种拦截器获取传入文件的大小?

【问题讨论】:

  • 随着内容的流式传输,它会随着时间的推移到达(因此在调用拦截器时信息不可用;您需要等待handler)。您是否希望在收到每个 PdfFile 时提出一个日志条目(在处理程序中而不是在拦截器中更容易做到这一点)或只是处理的总字节数(仅在流关闭时才知道;认为这需要一个 @987654322 @)。
  • 是的,你是对的,这实际上是我成功实现这一点的方式!

标签: go grpc middleware interceptor go-grpc-middleware


【解决方案1】:

所以,正如@Brits 上面提到的,实现我想要的工作方式是为流编写包装器。

这是一个示例:https://github.com/grpc-ecosystem/go-grpc-middleware/blob/master/validator/validator.go 我从这个 repo 中获取了以下代码,我希望我理解 apache2 许可证是正确的,并且复制此代码的一部分没有问题:

// StreamServerInterceptor returns a new streaming server interceptor that validates incoming messages.
//
// The stage at which invalid messages will be rejected with `InvalidArgument` varies based on the
// type of the RPC. For `ServerStream` (1:m) requests, it will happen before reaching any userspace
// handlers. For `ClientStream` (n:1) or `BidiStream` (n:m) RPCs, the messages will be rejected on
// calls to `stream.Recv()`.
func StreamServerInterceptor() grpc.StreamServerInterceptor {
    return func(srv interface{}, stream grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
        wrapper := &recvWrapper{stream}
        return handler(srv, wrapper)
    }
}

type recvWrapper struct {
    grpc.ServerStream
}

func (s *recvWrapper) RecvMsg(m interface{}) error {
    if err := s.ServerStream.RecvMsg(m); err != nil {
        return err
    }

    if err := validate(m); err != nil {
        return err
    }

    return nil
}

validate() 函数实际上将请求获取为interface{},您需要使用类型断言将此请求转换为您需要的类型为v := req.(type)

例如,我向FileForm 输入请求并能够检查其内容:

func logFileSize(req interface{}) error {
    m := req.(*pdfcompose.FileForm)
    println("SizeOfUpfile1: " + strconv.Itoa(int(binary.Size(m.Upfile1))))

    if m.Upfile2 != nil {
        println("SizeOfUpfile2: " + strconv.Itoa(int(binary.Size(m.Upfile3))))
    }
    if m.Upfile3 != nil {
        println("SizeOfUpfile2: " + strconv.Itoa(int(binary.Size(m.Upfile3))))
    }
    return nil
}

func (s *recvWrapper) RecvMsg(m interface{}) error {
    if err := s.ServerStream.RecvMsg(m); err != nil {
        return err
    }
    //z := m.(pdfcompose.FileForm)
    if err := logFileSize(m); err != nil {
        return err
    }

    return nil
}

func StreamServerInterceptor() grpc.StreamServerInterceptor {
    return func(srv interface{}, stream grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
        wrapper := &recvWrapper{stream}
        return handler(srv, wrapper)
    }
}

【讨论】:

    猜你喜欢
    • 2019-11-24
    • 2016-11-13
    • 1970-01-01
    • 1970-01-01
    • 2016-04-26
    • 1970-01-01
    • 2019-02-11
    • 2021-05-11
    • 2022-12-03
    相关资源
    最近更新 更多