ããã«ã¡ã¯ïŒä»å¹Žã® 4 æããã¹ããŒããã£ã³ãã«å
¥ç€Ÿããåªä»æ°åãšã³ãžãã¢ç ä¿®æéäžã®äžç°ã§ããæ¬èšäºã¯ãã€ã³ã¿ãŒãã§ãŒã¹å®çŸ©ã®æ©ã¿ã解決ããããã« gRPCãProtocol Buffers ã調æ»ããŠã¿ãïŒãšããå
容ã®ãšã³ããªã§ãã èæ¯ gRPC ãšã¯ Protocol Buffers ãšã¯ 4 ã€ã®éä¿¡æ¹åŒã詊ããŠã¿ã å®è£
æºå ã€ã³ã¿ãŒãã§ãŒã¹å®çŸ© ã³ã³ãã€ã« ãµãŒããŒãšã¯ã©ã€ã¢ã³ãã®å®è£
UnaryCall ClientStreamingCall ServerStreamingCall BidirectionalStreamingCall ããã¥ã¡ã³ãçæ åŠã³ ãŸãšã èæ¯ æ°åãšã³ãžãã¢ç ä¿®ã§ã¯ãåæã®ã¡ã³ããŒãš 2 人㧠Go (REST API) + React/TS æ§æã® SPA ãäœã£ãŠããŸãã ãã®ã¢ããªã®éçºã§ã¯ãServer - Client éã§ã€ã³ã¿ãŒãã§ãŒã¹ã®å®çŸ©ãäžå
åã§ããŠããããããããã§ãªã¯ãšã¹ã/ã¬ã¹ãã³ã¹ã®åãå®çŸ©ããŠããããã«ãäžæ¹ã«ä¿®æ£ãå
¥ã£ãã¿ã€ãã³ã°ã§ããäžæ¹ãä¿®æ£ããå¿
èŠããããé¢åã ãªããšæããŠããŸããã ãŸããã€ã³ã¿ãŒãã§ãŒã¹ã®å®çŸ©å
容ãããã¥ã¡ã³ããªã©ã§ç®¡çããŠããªããããéçºã¡ã³ããŒéã§ã®èªèã«ãããçããŠããŸãããšããããŸããã ãã®èŸºãã®è§£æ±ºçãèããŠããéã«ãgRPC ã Protocol Buffers ã«ã€ããŠç¥ããèå³ãæã£ãã®ã§èª¿æ»ããŠã¿ãŸããã gRPC ãšã¯ gRPC 㯠Google ãéçºãã RPC ã®ãã¬ãŒã ã¯ãŒã¯ã§ãã ããŒã¿ã®ã·ãªã¢ã©ã€ãºãšã€ã³ã¿ãŒãã§ãŒã¹ã®å®çŸ©ã« Protocol Buffers ãçšããããã°ã©ãã³ã°èšèªã«äŸåããªãå®è£
ã§é«éãªéä¿¡ãå¯èœã«ãªããšããç¹åŸŽããããŸãã Protocol Buffers ãšã¯ æ§é åããŒã¿ããããã¯ãŒã¯çµç±ã§éä¿¡ã§ãã圢ãžã·ãªã¢ã©ã€ãºãã圹å²ãšãæ§é åããŒã¿å®çŸ©çšã® IDL(Interface Depription Language)ãšããŠã®åœ¹å²ãæ©èœãšããŠåããŠããŸãã å
¬åŒã§ã¯ä»¥äžã®ããã«å®çŸ©ãããŠããŸãã ãããã³ã«ãããã¡ã¯ãæ§é åãããããŒã¿ãã·ãªã¢ã«åããããã®ãèšèªããã©ãããã©ãŒã ã«äŸåããªããGoogle ã®æ¡åŒµå¯èœãªã¡ã«ããºã ã§ãã åŒçšïŒ https://developers.google.com/protocol-buffers/ 4 ã€ã®éä¿¡æ¹åŒã詊ããŠã¿ã â»æ¬èšäºã§æ±ã£ãŠããã³ãŒã㯠ã³ã³ ã«çœ®ããŠãããŸãã gRPC ã§ã¯ Stream ãå©çšããŠã1 ãªã¯ãšã¹ãå
ã§ã¯ã©ã€ã¢ã³ããšãµãŒããŒéã§ã®ã¡ãã»ãŒãžã®ãããšããè€æ°åè¡ãããšãã§ããŸãã ããã«ãããgRPC ã§ã¯ä»¥äžã® 4 ã€ã®æ¹åŒã§éä¿¡ãã§ããŸãã Unary RPC (ã·ã³ãã« RPC) Server streaming RPC (ãµãŒããŒãµã€ãã¹ããªãŒãã³ã° RPC) Client streaming RPC (ã¯ã©ã€ã¢ã³ããµã€ãã¹ããªãŒãã³ã° RPC) Bidirectional streaming RPC (åæ¹åã¹ããªãŒãã³ã° RPC) åŒçšïŒ https://grpc.io/docs/what-is-grpc/core-concepts/#rpc-life-cycle ããããã«ããªã¯ãšã¹ããšã¬ã¹ãã³ã¹ã®é¢ä¿ã以äžã®ããã«å€åããŸãã Unary RPC - ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹ = 1:1 Server streaming RPC - ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹ = 1:N Client streaming RPC - ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹ = N:1 Bidirectional streaming RPC - ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹ = N:N ä»åã¯äžèšã®åéä¿¡æ¹åŒã Go ã§å®è£
ããŠã¿ãŸãã å®è£
æºå åãã«ã以äžãã€ã³ã¹ããŒã«ããŸãã Protocol Buffers v3 grpc-go( https://pkg.go.dev/google.golang.org/grpc ) protoc-gen-go( https://pkg.go.dev/github.com/golang/protobuf/protoc-gen-go ) > brew install protobuf > go get -u google.golang.org/grpc > go get -u github.com/golang/protobuf/protoc-gen-go ã€ã³ã¿ãŒãã§ãŒã¹å®çŸ© 次ã«ã.proto ãã¡ã€ã«ã«ã€ã³ã¿ãŒãã§ãŒã¹ã®å®çŸ©ãæžããŠãããŸãã ä»å㯠4 ã€ã®éä¿¡æ¹åŒã詊ããããservice Call ã«å¯ŸããŠåéä¿¡æ¹åŒã® service ã¡ãœãããå®çŸ©ããŠããŸãã ./proto/call.proto syntax = "proto3" ; package call; option go_package = "gen/pb" ; service Call { // Unary Request:Response = 1:1 rpc UnaryCall (CallRequest) returns (CallResponse) {} // ClientStreaming Request:Response = N:1 rpc ClientStreamingCall (stream CallRequest) returns (CallResponse) {} // ServerStreaming Request:Response = 1:N rpc ServerStreamingCall (ServerStreamingCallRequest) returns (stream CallResponse) {} // BidirectionalStreaming Request:Response = N:N rpc BidirectionalStreamingCall (stream CallRequest) returns (stream BidirectionalStreamingResponse) {} } message CallRequest { string name = 1 ; } message CallResponse { string message = 1 ; } message ServerStreamingCallRequest { string name = 1 ; uint32 responseCnt = 2 ; } message BidirectionalStreamingResponse { map< string , uint32 > callCounter = 1 ; } çŽæçãªèšæ³ã§å®çŸ©ã§ããŠåããããããªãšããå°è±¡ã§ããã ã³ã³ãã€ã« å®çŸ©ãã.proto ãã¡ã€ã«ã protoc ã§ã³ã³ãã€ã«ããŸãã > protoc --proto_path ./proto --go_out=plugins=grpc: ${APP_ROOT} call.proto ./gen/pb 以äžã« call.pb.go ãã¡ã€ã«ãèªåçæãããŸããã çæããã call.pb.go ã«ã¯ä»¥äžã®ãããªå
容ãå«ãŸããŠããŸãã å®çŸ©ãã message ã«å¯Ÿå¿ããæ§é äœ type CallRequest struct { state protoimpl.MessageState sizeCache protoimpl.SizeCache unknownFields protoimpl.UnknownFields Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` } service ã¡ãœãããåŒã³åºã gRPC ã¯ã©ã€ã¢ã³ã type CallClient interface { // Unary Request:Response = 1:1 UnaryCall(ctx context.Context, in *CallRequest, opts ...grpc.CallOption) (*CallResponse, error ) // ClientStreaming Request:Response = N:1 ClientStreamingCall(ctx context.Context, opts ...grpc.CallOption) (Call_ClientStreamingCallClient, error ) // ServerStreaming Request:Response = 1:N ServerStreamingCall(ctx context.Context, in *ServerStreamingCallRequest, opts ...grpc.CallOption) (Call_ServerStreamingCallClient, error ) // BidirectionalStreaming Request:Response = N:N BidirectionalStreamingCall(ctx context.Context, opts ...grpc.CallOption) (Call_BidirectionalStreamingCallClient, error ) } type callClient struct { cc grpc.ClientConnInterface } func NewCallClient(cc grpc.ClientConnInterface) CallClient { return &callClient{cc} } func (c *callClient) UnaryCall(ctx context.Context, in *CallRequest, opts ...grpc.CallOption) (*CallResponse, error ) { out := new (CallResponse) err := c.cc.Invoke(ctx, "/call.Call/UnaryCall" , in, out, opts...) if err != nil { return nil , err } return out, nil } service ã¡ãœãããå®è£
ãã gRPC ãµãŒããŒã®ã€ã³ã¿ãŒãã§ãŒã¹ type CallServer interface { // Unary Request:Response = 1:1 UnaryCall(context.Context, *CallRequest) (*CallResponse, error ) // ClientStreaming Request:Response = N:1 ClientStreamingCall(Call_ClientStreamingCallServer) error // ServerStreaming Request:Response = 1:N ServerStreamingCall(*ServerStreamingCallRequest, Call_ServerStreamingCallServer) error // BidirectionalStreaming Request:Response = N:N BidirectionalStreamingCall(Call_BidirectionalStreamingCallServer) error } ãµãŒããŒãšã¯ã©ã€ã¢ã³ãã®å®è£
次ã«ãçæããã³ãŒããçšããŠããµãŒããŒãšã¯ã©ã€ã¢ã³ãã®å®è£
ãããŠãããŸãã å®è£
ã¯å
¬åŒã«çšæãããŠãã åéä¿¡æ¹åŒã®å®è£
äŸ ãåèã«é²ããŸããã å®éã«æžããã³ãŒãã¯ä»¥äžã«ãªããŸãã ./client/cmd/main.go package main import ( "context" "grpc-lesson/gen/pb" "io" "log" "time" "google.golang.org/grpc" ) const ( address = "localhost:50051" ) func runUnaryCall(c pb.CallClient, name string ) error { log.Println( "--- Unary ---" ) in := &pb.CallRequest{Name: name} ctx, cancel := context.WithTimeout(context.Background(), time.Second* 10 ) defer cancel() res, err := c.UnaryCall(ctx, in) if err != nil { return err } log.Printf( "response: %s \n " , res.GetMessage()) return nil } func runClientStreamingCall(c pb.CallClient, names [] string ) error { log.Println( "--- ClientStreaming ---" ) ctx, cancel := context.WithTimeout(context.Background(), time.Second* 10 ) defer cancel() stream, err := c.ClientStreamingCall(ctx) if err != nil { return err } for _, name := range names { in := &pb.CallRequest{Name: name} if err := stream.Send(in); err != nil { if err == io.EOF { break } return err } time.Sleep(time.Second) } res, err := stream.CloseAndRecv() if err != nil { return err } log.Printf( "response: %s" , res.GetMessage()) return nil } func runServerStreamingCall(c pb.CallClient, name string , responseCnt uint32 ) error { log.Println( "--- ServerStreaming ---" ) in := &pb.ServerStreamingCallRequest{Name: name, ResponseCnt: responseCnt} ctx, cancel := context.WithTimeout(context.Background(), time.Second* 30 ) defer cancel() stream, err := c.ServerStreamingCall(ctx, in) if err != nil { return err } for { res, err := stream.Recv() if err == io.EOF { break } if err != nil { return err } log.Printf( "response: %s" , res.GetMessage()) } return nil } func runBidirectionalStreamingCall(c pb.CallClient, names [] string ) error { log.Println( "--- BidirectionalStreaming ---" ) ctx, cancel := context.WithTimeout(context.Background(), time.Minute) defer cancel() stream, err := c.BidirectionalStreamingCall(ctx) if err != nil { return err } done := make ( chan struct {}) go recv(done, stream) if err := send(names, stream); err != nil { return err } <- done return nil } func send(names [] string , stream pb.Call_BidirectionalStreamingCallClient) error { for _, name := range names { in := &pb.CallRequest{Name:name} if err := stream.Send(in); err != nil { return err } } if err := stream.CloseSend(); err != nil { return err } return nil } func recv(done chan struct {}, stream pb.Call_BidirectionalStreamingCallClient) { for { res, err := stream.Recv() if err == io.EOF { close (done) return } if err != nil { log.Fatalln(err) } log.Printf( "response: %v \n " , res.CallCounter) } } func main() { conn, err := grpc.Dial(address, grpc.WithInsecure(), grpc.WithBlock()) if err != nil { log.Fatalf( "did not connect %v" , err) } defer conn.Close() c := pb.NewCallClient(conn) if err = runUnaryCall(c, "John" ); err != nil { log.Fatalln(err) } names := [] string { "John" , "Paul" , "George" , "Ringo" } if err = runClientStreamingCall(c, names); err != nil { log.Fatalln(err) } if err = runServerStreamingCall(c, "John" , 10 ); err != nil { log.Fatalln(err) } names = [] string { "John" , "Paul" , "John" , "George" , "Ringo" , "Paul" , "John" , "Paul" , "George" , "John" } if err = runBidirectionalStreamingCall(c, names); err != nil { log.Fatalln(err) } } ./server/cmd/main.go package main import ( "context" "errors" "fmt" "grpc-lesson/gen/pb" "io" "log" "net" "time" "google.golang.org/grpc" ) const port = ":50051" type CallServer struct { pb.UnimplementedCallServer } func (s *CallServer) UnaryCall(ctx context.Context, in *pb.CallRequest) (*pb.CallResponse, error ) { log.Println( "--- Unary ---" ) log.Printf( "request: %s \n " , in.GetName()) resp := &pb.CallResponse{} resp.Message = fmt.Sprintf( "Hello. I'm %s" , in.GetName()) return resp, nil } func (s *CallServer) ClientStreamingCall(stream pb.Call_ClientStreamingCallServer) error { log.Println( "--- ClientStreaming ---" ) message := "Hello. We're" for { in, err := stream.Recv() if err == io.EOF { return stream.SendAndClose(&pb.CallResponse{Message: message}) } if err != nil { return err } log.Printf( "request: %s \n " , in.GetName()) message = fmt.Sprintf( "%s %s" , message, in.GetName()) } } func (s *CallServer) ServerStreamingCall(in *pb.ServerStreamingCallRequest, stream pb.Call_ServerStreamingCallServer) error { log.Println( "--- ClientStreaming ---" ) log.Printf( "request: %s \n " , in.GetName()) var message string for i := uint32 ( 1 ); i <= in.ResponseCnt; i++ { if i <= 5 { message = fmt.Sprintf( "Hello. I'm %s" , in.GetName()) } else { message = fmt.Sprintf( "I'm so tired. (%s)" , in.GetName()) } if err := stream.Send(&pb.CallResponse{Message: message}); err != nil { return err } time.Sleep(time.Second) } return nil } func (s *CallServer) BidirectionalStreamingCall(stream pb.Call_BidirectionalStreamingCallServer) error { log.Println( "--- BidirectionalStreaming ---" ) counter := make ( map [ string ] uint32 ) for { in, err := stream.Recv() if err == io.EOF { return nil } if err != nil { return err } log.Printf( "request: %s \n " , in.GetName()) counter[in.Name] += 1 res := &pb.BidirectionalStreamingResponse{CallCounter: counter} if err := stream.Send(res); err != nil { return err } } } func main() { fmt.Printf( "server is listening on port%s \n " , port) if err := set(); err != nil { log.Fatalln(err.Error()) } } func set() error { lis, err := net.Listen( "tcp" , port) if err != nil { log.Fatalln(err) } s := grpc.NewServer() pb.RegisterCallServer(s, &CallServer{}) if err := s.Serve(lis); err != nil { return errors.New( "serve is failed" ) } return nil } åéä¿¡æ¹åŒã®å®è£
ãèŠãŠãããŸãã UnaryCall UnaryCall ã¯ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹=1:1 ã®éä¿¡ã§ãã ãœãŒã¹ã³ãŒã ./proto/call.proto service Call { // Unary Request:Response = 1:1 rpc UnaryCall (CallRequest) returns (CallResponse) {} ... } message CallRequest { string name = 1 ; } message CallResponse { string message = 1 ; } ./client/cmd/main.go func runUnaryCall(c pb.CallClient, name string ) error { log.Println( "--- Unary ---" ) in := &pb.CallRequest{Name: name} ctx, cancel := context.WithTimeout(context.Background(), time.Second* 10 ) defer cancel() res, err := c.UnaryCall(ctx, in) if err != nil { return err } log.Printf( "response: %s \n " , res.GetMessage()) return nil } func main() { ... c := pb.NewCallClient(conn) if err = runUnaryCall(c, "John" ); err != nil { log.Fatalln(err) } ... } UnaryCall ã®ã¯ã©ã€ã¢ã³ãã§ã¯ãCallRequest ããµãŒããŒãžæã CallResponse ãåãåãã¡ãã»ãŒãžã log ã«åããŠããŸãã ./server/cmd/main.go func (s *CallServer) UnaryCall(ctx context.Context, in *pb.CallRequest) (*pb.CallResponse, error ) { log.Println( "--- Unary ---" ) log.Printf( "request: %s \n " , in.GetName()) resp := &pb.CallResponse{} resp.Message = fmt.Sprintf( "Hello. I'm %s" , in.GetName()) return resp, nil } UnaryCall ã®ãµãŒããŒã§ã¯ãCallRequest ãåãåã£ãŠ log ãžåããåãåã£ã CallRequest.Name ãããšã« CallResponse ãçæããŠã¯ã©ã€ã¢ã³ããè¿ããŠããŸãã å®è¡çµæ ãµãŒã㌠2021/06/08 00:57:37 --- Unary --- 2021/06/08 00:57:37 request: John ã¯ã©ã€ã¢ã³ã 2021/06/08 00:57:37 --- Unary --- 2021/06/08 00:57:37 response: Hello. I'm John ClientStreamingCall ClientStreamingCall ã¯ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹=N:1 ã®éä¿¡ã§ãã ãœãŒã¹ã³ãŒã ./proto/call.proto service Call { ... // ClientStreaming Request:Response = N:1 rpc ClientStreamingCall (stream CallRequest) returns (CallResponse) {} ... } message CallRequest { string name = 1 ; } message CallResponse { string message = 1 ; } ./client/cmd/main.go func runClientStreamingCall(c pb.CallClient, names [] string ) error { log.Println( "--- ClientStreaming ---" ) ctx, cancel := context.WithTimeout(context.Background(), time.Second* 10 ) defer cancel() stream, err := c.ClientStreamingCall(ctx) if err != nil { return err } for _, name := range names { in := &pb.CallRequest{Name: name} if err := stream.Send(in); err != nil { if err == io.EOF { break } return err } time.Sleep(time.Second) } res, err := stream.CloseAndRecv() if err != nil { return err } log.Printf( "response: %s" , res.GetMessage()) return nil } func main() { ... c := pb.NewCallClient(conn) ... names := [] string { "John" , "Paul" , "George" , "Ringo" } if err = runClientStreamingCall(c, names); err != nil { log.Fatalln(err) } ... } ClientStreamingCall ã®ã¯ã©ã€ã¢ã³ãã§ã¯ã1 ç§ããã« CallRequest ããµãŒããŒãžæããå
šãŠæãçµãã£ãã 1 ã€ã®ã¬ã¹ãã³ã¹ãåãåããåãåã£ãã¡ãã»ãŒãžã log ã«åããŠããŸãã ./server/cmd/main.go func (s *CallServer) ClientStreamingCall(stream pb.Call_ClientStreamingCallServer) error { log.Println( "--- ClientStreaming ---" ) message := "Hello. We're" for { in, err := stream.Recv() if err == io.EOF { return stream.SendAndClose(&pb.CallResponse{Message: message}) } if err != nil { return err } log.Printf( "request: %s \n " , in.GetName()) message = fmt.Sprintf( "%s %s" , message, in.GetName()) } } ClientStreamingCall ã®ãµãŒããŒã§ã¯ãCallRequest ãè€æ°ååãåã£ãŠ log ãžåããåãåã£ãè€æ°ã®ãªã¯ãšã¹ããããšã« 1 ã€ã® CallResponse ãçæããŠè¿ããŠããŸãã å®è¡çµæ ãµãŒã㌠2021/06/08 00:57:37 --- ClientStreaming --- 2021/06/08 00:57:37 request: John 2021/06/08 00:57:38 request: Paul 2021/06/08 00:57:39 request: George 2021/06/08 00:57:40 request: Ringo ã¯ã©ã€ã¢ã³ã 2021/06/08 00:57:37 --- ClientStreaming --- 2021/06/08 00:57:41 response: Hello. We're John Paul George Ringo ServerStreamingCall ServerStreamingCall ã¯ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹=1:N ã®éä¿¡ã§ãã ãœãŒã¹ã³ãŒã ./proto/call.proto service Call { ... // ServerStreaming Request:Response = 1:N rpc ServerStreamingCall (ServerStreamingCallRequest) returns (stream CallResponse) {} ... } message ServerStreamingCallRequest { string name = 1 ; uint32 responseCnt = 2 ; } message CallResponse { string message = 1 ; } ./client/cmd/main.go func runServerStreamingCall(c pb.CallClient, name string , responseCnt uint32 ) error { log.Println( "--- ServerStreaming ---" ) in := &pb.ServerStreamingCallRequest{Name: name, ResponseCnt: responseCnt} ctx, cancel := context.WithTimeout(context.Background(), time.Second* 30 ) defer cancel() stream, err := c.ServerStreamingCall(ctx, in) if err != nil { return err } for { res, err := stream.Recv() if err == io.EOF { break } if err != nil { return err } log.Printf( "response: %s" , res.GetMessage()) } return nil } func main() { ... c := pb.NewCallClient(conn) ... if err = runServerStreamingCall(c, "John" , 10 ); err != nil { log.Fatalln(err) } ... } ServerStreamingCall ã®ã¯ã©ã€ã¢ã³ãã§ã¯ãName ãš ResponseCnt ãšãããã£ãŒã«ããæã€ ServerStreamingCallRequest ããµãŒããŒãžæããè€æ°åã® CallResponse ãåãåããåãåã£ãããããã®ã¡ãã»ãŒãžã log ã«åããŠããŸãã ./server/cmd/main.go func (s *CallServer) ServerStreamingCall(in *pb.ServerStreamingCallRequest, stream pb.Call_ServerStreamingCallServer) error { log.Println( "--- ClientStreaming ---" ) log.Printf( "request: %s \n " , in.GetName()) var message string for i := uint32 ( 1 ); i <= in.ResponseCnt; i++ { if i <= 5 { message = fmt.Sprintf( "Hello. I'm %s" , in.GetName()) } else { message = fmt.Sprintf( "I'm so tired. (%s)" , in.GetName()) } if err := stream.Send(&pb.CallResponse{Message: message}); err != nil { return err } time.Sleep(time.Second) } return nil } ServerStreamingCall ã®ãµãŒããŒã§ã¯ãServerStreamingCallRequest ãåãåã£ãŠ Name ã log ãžåããåãåã£ããªã¯ãšã¹ãã® ResponseCnt ãã£ãŒã«ãã§æå®ãããåæ°ã ã CallResponse ãçæããŠè¿ããŠããŸãã å®è¡çµæ ãµãŒã㌠2021/06/08 01:29:03 --- ServerStreaming --- 2021/06/08 01:29:03 request: John ã¯ã©ã€ã¢ã³ã 2021/06/08 01:29:03 --- ServerStreaming --- 2021/06/08 01:29:03 response: Hello. I'm John 2021/06/08 01:29:04 response: Hello. I'm John 2021/06/08 01:29:05 response: Hello. I'm John 2021/06/08 01:29:06 response: Hello. I'm John 2021/06/08 01:29:07 response: Hello. I'm John 2021/06/08 01:29:08 response: I'm so tired. (John) 2021/06/08 01:29:09 response: I'm so tired. (John) 2021/06/08 01:29:10 response: I'm so tired. (John) 2021/06/08 01:29:11 response: I'm so tired. (John) BidirectionalStreamingCall BidirectionalStreamingCall ã¯ãªã¯ãšã¹ã:ã¬ã¹ãã³ã¹=N:N ã®åæ¹åã®éä¿¡ã§ãã ãœãŒã¹ã³ãŒã ./proto/call.proto service Call { ... // BidirectionalStreaming Request:Response = N:N rpc BidirectionalStreamingCall (stream CallRequest) returns (stream BidirectionalStreamingResponse) {} } message ServerStreamingCallRequest { string name = 1 ; uint32 responseCnt = 2 ; } message BidirectionalStreamingResponse { map< string , uint32 > callCounter = 1 ; } ./client/cmd/main.go func runBidirectionalStreamingCall(c pb.CallClient, names [] string ) error { log.Println( "--- BidirectionalStreaming ---" ) ctx, cancel := context.WithTimeout(context.Background(), time.Minute) defer cancel() stream, err := c.BidirectionalStreamingCall(ctx) if err != nil { return err } done := make ( chan struct {}) go recv(done, stream) if err := send(names, stream); err != nil { return err } <- done return nil } func send(names [] string , stream pb.Call_BidirectionalStreamingCallClient) error { for _, name := range names { in := &pb.CallRequest{Name:name} if err := stream.Send(in); err != nil { return err } } if err := stream.CloseSend(); err != nil { return err } return nil } func recv(done chan struct {}, stream pb.Call_BidirectionalStreamingCallClient) { for { res, err := stream.Recv() if err == io.EOF { close (done) return } if err != nil { log.Fatalln(err) } log.Printf( "response: %v \n " , res.CallCounter) } } func main() { ... c := pb.NewCallClient(conn) ... names = [] string { "John" , "Paul" , "John" , "George" , "Ringo" , "Paul" , "John" , "Paul" , "George" , "John" } if err = runBidirectionalStreamingCall(c, names); err != nil { log.Fatalln(err) } ... } BidirectionalStreamingCall ã®ã¯ã©ã€ã¢ã³ãã§ã¯ãCallRequest ãè€æ°åãµãŒããŒãžæããBidirectionalCallResponse ãè€æ°ååãåããã¬ã¹ãã³ã¹ã® CallCounter ããããã log ã«åããŠããŸãã ./server/cmd/main.go func (s *CallServer) BidirectionalStreamingCall(stream pb.Call_BidirectionalStreamingCallServer) error { log.Println( "--- BidirectionalStreaming ---" ) counter := make ( map [ string ] uint32 ) for { in, err := stream.Recv() if err == io.EOF { return nil } if err != nil { return err } log.Printf( "request: %s \n " , in.GetName()) counter[in.Name] += 1 res := &pb.BidirectionalStreamingResponse{CallCounter: counter} if err := stream.Send(res); err != nil { return err } } } BidirectionalStreamingCall ã®ãµãŒããŒã§ã¯ãCallRequest ãè€æ°ååãã€ããåãªã¯ãšã¹ããåãåã床㫠BidirectionalStreamingCallResponse ã® CallCounter ã«ååãåŒã°ããæ°ã CountUp ããŠã¬ã¹ãã³ã¹ãšããŠè¿åŽããŠããŸãã å®è¡çµæ ãµãŒã㌠2021/06/08 01:29:13 --- BidirectionalStreaming --- 2021/06/08 01:29:13 request: John 2021/06/08 01:29:13 request: Paul 2021/06/08 01:29:13 request: John 2021/06/08 01:29:13 request: George 2021/06/08 01:29:13 request: Ringo 2021/06/08 01:29:13 request: Paul 2021/06/08 01:29:13 request: John 2021/06/08 01:29:13 request: Paul 2021/06/08 01:29:13 request: George 2021/06/08 01:29:13 request: John ã¯ã©ã€ã¢ã³ã 2021/06/08 01:29:13 --- BidirectionalStreaming --- 2021/06/08 01:29:13 response: map[John:1] 2021/06/08 01:29:13 response: map[John:1 Paul:1] 2021/06/08 01:29:13 response: map[John:2 Paul:1] 2021/06/08 01:29:13 response: map[George:1 John:2 Paul:1] 2021/06/08 01:29:13 response: map[George:1 John:2 Paul:1 Ringo:1] 2021/06/08 01:29:13 response: map[George:1 John:2 Paul:2 Ringo:1] 2021/06/08 01:29:13 response: map[George:1 John:3 Paul:2 Ringo:1] 2021/06/08 01:29:13 response: map[George:1 John:3 Paul:3 Ringo:1] 2021/06/08 01:29:13 response: map[George:2 John:3 Paul:3 Ringo:1] 2021/06/08 01:29:13 response: map[George:2 John:4 Paul:3 Ringo:1] ããã¥ã¡ã³ãçæ æåŸã«ã€ã³ã¿ãŒãã§ãŒã¹å®çŸ©ã®ããã¥ã¡ã³ããäœæããŸãã protoc ã®ããã¥ã¡ã³ãçæçšã®ãã©ã°ã€ã³ãã€ã³ã¹ããŒã«ããŸãã > go get -u github.com/pseudomuto/protoc-gen-doc/cmd/protoc-gen-doc ãã® protoc-gen-doc ã§ã¯ã.proto ãã¡ã€ã«ãã HTML,JSON,DocBook,Markdown ã®ããããã®åœ¢åŒã§ããã¥ã¡ã³ããçæããããšãã§ããŸãã ä»å㯠HTML 圢åŒã§ããã¥ã¡ã³ããçæããŸãã > protoc --doc_out=./proto/doc --doc_opt=html,index.html ./proto/*.proto ./proto/doc 以äžã« index.html ãšããããã¥ã¡ã³ããã¡ã€ã«ãçæãããŸããã index.html ã³ãã³ã 1 çºã§æŽã£ãããã¥ã¡ã³ããäœãããŠãšãŠã䟿å©ã§ããã åŠã³ 4 ã€ã®éä¿¡æ¹åŒã®ãããããprotobuf ã§çæãããæ
å ±ãããšã«ããã»ã©èŠåŽããã«å®è£
ããããšãã§ããŸããã ä»åã¯ã詊ããŠã¿ããèšäºãšããããšã§åéä¿¡æ¹åŒã®ç¹åŸŽãçãããã³ãŒãã¯æžããŠããŸããããå
¬åŒã®äŸã®ããã« ServerStreaming ã§ããã°ãµãŒããŒããã¯ã©ã€ã¢ã³ããžã® Push éç¥æ©èœãBidirectionalStreaming ã§ããã°ãµãŒããŒãä»ããŠã®ãã£ããæ©èœãªã©ãå®çŸãããäºã«å¿ããŠéä¿¡æ¹åŒã䜿ãåãããšåéä¿¡æ¹åŒã®è¯ãããã宿ã§ãããã ãªãšæããŸããã ãŸããæ°åãšã³ãžãã¢ç ä¿®ã§äœæããŠãã REST API ã®ã¢ããªã±ãŒã·ã§ã³ã§ããprotobuf ã§ ãªã¯ãšã¹ããšã¬ã¹ãã³ã¹ã® message ã®ã¿å®çŸ©ãããããªäœ¿ãæ¹ã§å®çŸ©ã®äžå
åã«åãçµãããã ãªãšæããŸãããããã¥ã¡ã³ãã®çæãããªã楜ã§è¯ãã£ãã§ãã ãŸãšã ãããã§ããã§ãããããä»å㯠gRPC ã® 4 ã€ã®éä¿¡æ¹åŒã Go ã§è©ŠããŠã¿ãŸããã åã¯ãããŸã§ REST ã§ã®éçºçµéšãã»ãšãã©ãªã®ã§ãä»ã®æ¹æ³ã§ã® API èšèšã¯æ°é®®ã§æ¥œããã£ãã§ããããªãã·ã³ãã«ã«éçºã§ãããã ã£ãã®ã§ãä»åŸæ¥åã§ãæ©äŒãèŠãŠå©çšããŠãããã°ãšæããŸãã æåŸãŸã§ãèªã¿ãã ããããããšãããããŸããïŒ