Chuyển đến nội dung chính

レッスン 14: gRPC とプロトコル バッファー

プロトコル バッファーのスキーマ設計、コード生成。単項、サーバー ストリーミング、クライアント ストリーミング、双方向ストリーミング。 gRPC インターセプター、エラー処理、期限。 REST 互換性のための gRPC ゲートウェイ。

💻 プログラミング — レッスン 14 レッスン 14: gRPC とプロトコル バッファー

Golang: 基本から高度まで

パート 4: 高度な機能

xdev.asia

1. プロトコルバッファ

1.1.セットアップ

# Install protoc compiler
brew install protobuf

# Install Go plugins
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

# Go dependencies
go get google.golang.org/grpc
go get google.golang.org/protobuf

1.2.プロトファイル

// proto/user/v1/user.proto
syntax = "proto3";

package user.v1;

option go_package = "myapp/gen/user/v1;userv1";

import "google/protobuf/timestamp.proto";
import "google/protobuf/empty.proto";

// Messages
message User {
    uint64 id = 1;
    string name = 2;
    string email = 3;
    Role role = 4;
    bool is_active = 5;
    google.protobuf.Timestamp created_at = 6;
}

enum Role {
    ROLE_UNSPECIFIED = 0;
    ROLE_USER = 1;
    ROLE_EDITOR = 2;
    ROLE_ADMIN = 3;
}

message CreateUserRequest {
    string name = 1;
    string email = 2;
    string password = 3;
}

message GetUserRequest {
    uint64 id = 1;
}

message ListUsersRequest {
    int32 page = 1;
    int32 page_size = 2;
    string search = 3;
}

message ListUsersResponse {
    repeated User users = 1;
    int32 total = 2;
    int32 page = 3;
}

message UpdateUserRequest {
    uint64 id = 1;
    optional string name = 2;
    optional string email = 3;
    optional Role role = 4;
}

// Service
service UserService {
    rpc CreateUser(CreateUserRequest) returns (User);
    rpc GetUser(GetUserRequest) returns (User);
    rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);
    rpc UpdateUser(UpdateUserRequest) returns (User);
    rpc DeleteUser(GetUserRequest) returns (google.protobuf.Empty);
    
    // Server streaming
    rpc WatchUsers(google.protobuf.Empty) returns (stream User);
    
    // Bidirectional streaming
    rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message ChatMessage {
    string sender = 1;
    string content = 2;
    google.protobuf.Timestamp sent_at = 3;
}

1.3.コード生成

# Generate Go code
protoc --go_out=. --go_opt=paths=source_relative \
       --go-grpc_out=. --go-grpc_opt=paths=source_relative \
       proto/user/v1/user.proto

# Makefile
.PHONY: proto
proto:
	protoc --go_out=. --go_opt=paths=source_relative \
	       --go-grpc_out=. --go-grpc_opt=paths=source_relative \
	       proto/**/**/*.proto

2. gRPC サーバー

package grpcserver

import (
    "context"
    "fmt"
    
    userv1 "myapp/gen/user/v1"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    "google.golang.org/protobuf/types/known/emptypb"
    "google.golang.org/protobuf/types/known/timestamppb"
)

type UserServer struct {
    userv1.UnimplementedUserServiceServer
    repo repository.UserRepository
}

func NewUserServer(repo repository.UserRepository) *UserServer {
    return &UserServer{repo: repo}
}

func (s *UserServer) CreateUser(ctx context.Context, req *userv1.CreateUserRequest) (*userv1.User, error) {
    // Validate input
    if req.Name == "" {
        return nil, status.Error(codes.InvalidArgument, "name is required")
    }
    if req.Email == "" {
        return nil, status.Error(codes.InvalidArgument, "email is required")
    }
    
    user, err := s.repo.Create(ctx, &model.User{
        Name:  req.Name,
        Email: req.Email,
    })
    if err != nil {
        return nil, status.Errorf(codes.Internal, "create user: %v", err)
    }
    
    return toProtoUser(user), nil
}

func (s *UserServer) GetUser(ctx context.Context, req *userv1.GetUserRequest) (*userv1.User, error) {
    user, err := s.repo.GetByID(ctx, uint(req.Id))
    if err != nil {
        return nil, status.Error(codes.NotFound, "user not found")
    }
    
    return toProtoUser(user), nil
}

func (s *UserServer) ListUsers(ctx context.Context, req *userv1.ListUsersRequest) (*userv1.ListUsersResponse, error) {
    page := int(req.Page)
    if page < 1 {
        page = 1
    }
    pageSize := int(req.PageSize)
    if pageSize < 1 || pageSize > 100 {
        pageSize = 20
    }
    
    users, total, err := s.repo.List(ctx, page, pageSize, req.Search)
    if err != nil {
        return nil, status.Errorf(codes.Internal, "list users: %v", err)
    }
    
    protoUsers := make([]*userv1.User, len(users))
    for i, u := range users {
        protoUsers[i] = toProtoUser(&u)
    }
    
    return &userv1.ListUsersResponse{
        Users: protoUsers,
        Total: int32(total),
        Page:  int32(page),
    }, nil
}

func (s *UserServer) DeleteUser(ctx context.Context, req *userv1.GetUserRequest) (*emptypb.Empty, error) {
    if err := s.repo.Delete(ctx, uint(req.Id)); err != nil {
        return nil, status.Error(codes.NotFound, "user not found")
    }
    return &emptypb.Empty{}, nil
}

// Server streaming
func (s *UserServer) WatchUsers(req *emptypb.Empty, stream userv1.UserService_WatchUsersServer) error {
    // Subscribe to user changes
    ch := s.repo.Subscribe()
    defer s.repo.Unsubscribe(ch)
    
    for {
        select {
        case <-stream.Context().Done():
            return nil
        case user := <-ch:
            if err := stream.Send(toProtoUser(user)); err != nil {
                return err
            }
        }
    }
}

func toProtoUser(u *model.User) *userv1.User {
    return &userv1.User{
        Id:        uint64(u.ID),
        Name:      u.Name,
        Email:     u.Email,
        IsActive:  u.IsActive,
        CreatedAt: timestamppb.New(u.CreatedAt),
    }
}

3. サーバーのセットアップ

package main

import (
    "log"
    "net"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/reflection"
    
    userv1 "myapp/gen/user/v1"
)

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("Failed to listen: %v", err)
    }
    
    server := grpc.NewServer(
        grpc.ChainUnaryInterceptor(
            LoggingInterceptor,
            AuthInterceptor,
            RecoveryInterceptor,
        ),
        grpc.ChainStreamInterceptor(
            StreamLoggingInterceptor,
        ),
    )
    
    // Register services
    userv1.RegisterUserServiceServer(server, NewUserServer(userRepo))
    
    // Enable reflection for debugging (grpcurl, grpcui)
    reflection.Register(server)
    
    log.Printf("gRPC server listening on :50051")
    if err := server.Serve(lis); err != nil {
        log.Fatalf("Failed to serve: %v", err)
    }
}

4. インターセプター (ミドルウェア)

import (
    "context"
    "log"
    "time"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/metadata"
    "google.golang.org/grpc/status"
)

// Logging interceptor
func LoggingInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (interface{}, error) {
    start := time.Now()
    
    resp, err := handler(ctx, req)
    
    duration := time.Since(start)
    statusCode := codes.OK
    if err != nil {
        statusCode = status.Code(err)
    }
    
    log.Printf("gRPC %s | %s | %v | %v",
        info.FullMethod, statusCode, duration, err)
    
    return resp, err
}

// Auth interceptor
func AuthInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (interface{}, error) {
    // Skip auth for certain methods
    publicMethods := map[string]bool{
        "/user.v1.UserService/CreateUser": true,
    }
    if publicMethods[info.FullMethod] {
        return handler(ctx, req)
    }
    
    // Extract token from metadata
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok {
        return nil, status.Error(codes.Unauthenticated, "missing metadata")
    }
    
    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return nil, status.Error(codes.Unauthenticated, "missing token")
    }
    
    // Validate token...
    claims, err := validateToken(tokens[0])
    if err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid token")
    }
    
    // Set user info in context
    ctx = context.WithValue(ctx, "user_id", claims.UserID)
    
    return handler(ctx, req)
}

// Recovery interceptor
func RecoveryInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (resp interface{}, err error) {
    defer func() {
        if r := recover(); r != nil {
            log.Printf("Panic recovered in %s: %v", info.FullMethod, r)
            err = status.Errorf(codes.Internal, "internal error")
        }
    }()
    
    return handler(ctx, req)
}

5. gRPC クライアント

package main

import (
    "context"
    "log"
    "time"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials/insecure"
    
    userv1 "myapp/gen/user/v1"
)

func main() {
    conn, err := grpc.NewClient("localhost:50051",
        grpc.WithTransportCredentials(insecure.NewCredentials()),
    )
    if err != nil {
        log.Fatalf("Connect failed: %v", err)
    }
    defer conn.Close()
    
    client := userv1.NewUserServiceClient(conn)
    
    // Unary call with deadline
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()
    
    user, err := client.CreateUser(ctx, &userv1.CreateUserRequest{
        Name:  "Duy Tran",
        Email: "[email protected]",
    })
    if err != nil {
        st, ok := status.FromError(err)
        if ok {
            log.Printf("gRPC error: code=%s, msg=%s", st.Code(), st.Message())
        }
        return
    }
    
    log.Printf("Created user: %+v", user)
    
    // Server streaming
    stream, err := client.WatchUsers(ctx, &emptypb.Empty{})
    if err != nil {
        log.Fatal(err)
    }
    
    for {
        user, err := stream.Recv()
        if err != nil {
            break
        }
        log.Printf("User update: %+v", user)
    }
}

6. gRPC エラー処理

gRPC コードHTTPコードいつ使用するか
OK200成功
無効な引数400無効な入力
未認証401認証されていません
許可が拒否されました403権利なし
見つかりません404リソースが存在しません
すでに存在します409重複
内部500サーバーエラー
利用不可503サービスを一時停止しております
期限を過ぎました504タイムアウト
// Error with details
import "google.golang.org/genproto/googleapis/rpc/errdetails"

func detailedError() error {
    st := status.New(codes.InvalidArgument, "invalid input")
    
    st, _ = st.WithDetails(&errdetails.BadRequest{
        FieldViolations: []*errdetails.BadRequest_FieldViolation{
            {Field: "email", Description: "invalid email format"},
            {Field: "name", Description: "name is required"},
        },
    })
    
    return st.Err()
}

7. gRPC ゲートウェイ (REST 互換)

go install github.com/grpc-ecosystem/grpc-gateway/v2/protoc-gen-grpc-gateway@latest
go install github.com/grpc-ecosystem/grpc-gateway/v2/protoc-gen-openapiv2@latest
import "google/api/annotations.proto";

service UserService {
    rpc GetUser(GetUserRequest) returns (User) {
        option (google.api.http) = {
            get: "/api/v1/users/{id}"
        };
    }
    
    rpc CreateUser(CreateUserRequest) returns (User) {
        option (google.api.http) = {
            post: "/api/v1/users"
            body: "*"
        };
    }
    
    rpc ListUsers(ListUsersRequest) returns (ListUsersResponse) {
        option (google.api.http) = {
            get: "/api/v1/users"
        };
    }
}
// Gateway server - proxy HTTP → gRPC
func runGateway() error {
    ctx := context.Background()
    mux := runtime.NewServeMux()
    
    opts := []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
    
    err := userv1.RegisterUserServiceHandlerFromEndpoint(ctx, mux, "localhost:50051", opts)
    if err != nil {
        return err
    }
    
    log.Println("Gateway listening on :8080")
    return http.ListenAndServe(":8080", mux)
}

次の記事: メッセージキューとイベント駆動型アーキテクチャ — RabbitMQ、Kafka、およびイベント ソーシング。