gRPC案例

wen java案例 1

本文目录导读:

gRPC案例

  1. 项目结构
  2. 定义 proto 文件 (user.proto)
  3. 生成代码
  4. 服务端实现 (server/main.go)
  5. 客户端实现 (client/main.go)
  6. 初始化 go.mod
  7. 运行测试
  8. 输出示例

我将为您提供一个完整的gRPC案例,包括 proto 文件定义、服务端实现和客户端调用。

项目结构

grpc-demo/
├── proto/
│   └── user.proto
├── server/
│   └── main.go
├── client/
│   └── main.go
├── go.mod
└── go.sum

定义 proto 文件 (user.proto)

syntax = "proto3";
package user;
option go_package = "grpc-demo/proto/user";
// 用户服务定义
service UserService {
  // 获取用户信息(一元 RPC)
  rpc GetUser(GetUserRequest) returns (User) {}
  // 获取用户列表(服务端流式 RPC)
  rpc ListUsers(ListUsersRequest) returns (stream User) {}
  // 批量创建用户(客户端流式 RPC)
  rpc BatchCreateUsers(stream CreateUserRequest) returns (CreateUsersResponse) {}
  // 用户聊天(双向流式 RPC)
  rpc Chat(stream ChatMessage) returns (stream ChatMessage) {}
}
// 请求消息
message GetUserRequest {
  int32 id = 1;
}
message ListUsersRequest {
  int32 page = 1;
  int32 page_size = 2;
}
message CreateUserRequest {
  string name = 1;
  string email = 2;
}
// 响应消息
message User {
  int32 id = 1;
  string name = 2;
  string email = 3;
  string created_at = 4;
}
message CreateUsersResponse {
  int32 success_count = 1;
  repeated User users = 2;
}
// 聊天消息
message ChatMessage {
  string user_name = 1;
  string content = 2;
  string timestamp = 3;
}

生成代码

# 安装 protoc 和相关插件
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest
# 生成代码
protoc --go_out=. --go_opt=paths=source_relative \
  --go-grpc_out=. --go-grpc_opt=paths=source_relative \
  proto/user.proto

服务端实现 (server/main.go)

package main
import (
    "context"
    "fmt"
    "io"
    "log"
    "net"
    "time"
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    pb "grpc-demo/proto/user"
)
// 用户数据结构
type userData struct {
    ID   int32
    Name string
    Email string
}
// 服务端实现
type userService struct {
    pb.UnimplementedUserServiceServer
    users map[int32]*userData
    nextID int32
}
func newUserService() *userService {
    service := &userService{
        users: make(map[int32]*userData),
        nextID: 1,
    }
    // 添加一些初始数据
    service.users[1] = &userData{ID: 1, Name: "张三", Email: "zhangsan@example.com"}
    service.users[2] = &userData{ID: 2, Name: "李四", Email: "lisi@example.com"}
    service.nextID = 3
    return service
}
// 一元 RPC:获取用户信息
func (s *userService) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.User, error) {
    log.Printf("收到获取用户请求: ID=%d", req.Id)
    // 模拟延迟
    time.Sleep(100 * time.Millisecond)
    user, exists := s.users[req.Id]
    if !exists {
        return nil, status.Errorf(codes.NotFound, "用户不存在: ID=%d", req.Id)
    }
    return &pb.User{
        Id:        user.ID,
        Name:      user.Name,
        Email:     user.Email,
        CreatedAt: time.Now().Format("2006-01-02 15:04:05"),
    }, nil
}
// 服务端流式 RPC:获取用户列表
func (s *userService) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
    log.Printf("收到用户列表请求: page=%d, pageSize=%d", req.Page, req.PageSize)
    // 分页逻辑
    start := int((req.Page - 1) * req.PageSize)
    end := start + int(req.PageSize)
    users := make([]*userData, 0, len(s.users))
    for _, user := range s.users {
        users = append(users, user)
    }
    if start > len(users) {
        start = len(users)
    }
    if end > len(users) {
        end = len(users)
    }
    // 流式发送用户数据
    for _, user := range users[start:end] {
        if err := stream.Send(&pb.User{
            Id:        user.ID,
            Name:      user.Name,
            Email:     user.Email,
            CreatedAt: time.Now().Format("2006-01-02 15:04:05"),
        }); err != nil {
            return err
        }
        time.Sleep(200 * time.Millisecond) // 模拟流式传输
    }
    return nil
}
// 客户端流式 RPC:批量创建用户
func (s *userService) BatchCreateUsers(stream pb.UserService_BatchCreateUsersServer) error {
    log.Println("开始批量创建用户")
    var createdUsers []*userData
    for {
        req, err := stream.Recv()
        if err == io.EOF {
            // 所有用户创建完成
            response := &pb.CreateUsersResponse{
                SuccessCount: int32(len(createdUsers)),
            }
            for _, user := range createdUsers {
                response.Users = append(response.Users, &pb.User{
                    Id:        user.ID,
                    Name:      user.Name,
                    Email:     user.Email,
                    CreatedAt: time.Now().Format("2006-01-02 15:04:05"),
                })
            }
            return stream.SendAndClose(response)
        }
        if err != nil {
            return err
        }
        // 创建新用户
        newUser := &userData{
            ID:    s.nextID,
            Name:  req.Name,
            Email: req.Email,
        }
        s.users[newUser.ID] = newUser
        s.nextID++
        createdUsers = append(createdUsers, newUser)
        log.Printf("已创建用户: %s (%s)", newUser.Name, newUser.Email)
    }
}
// 双向流式 RPC:用户聊天
func (s *userService) Chat(stream pb.UserService_ChatServer) error {
    log.Println("开始聊天会话")
    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            return nil
        }
        if err != nil {
            return err
        }
        log.Printf("[%s]: %s", msg.UserName, msg.Content)
        // 模拟自动回复
        reply := &pb.ChatMessage{
            UserName:  "Server",
            Content:   fmt.Sprintf("收到消息: \"%s\"", msg.Content),
            Timestamp: time.Now().Format("2006-01-02 15:04:05"),
        }
        if err := stream.Send(reply); err != nil {
            return err
        }
    }
}
func main() {
    // 创建 gRPC 服务器
    server := grpc.NewServer()
    // 注册服务
    service := newUserService()
    pb.RegisterUserServiceServer(server, service)
    // 监听端口
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("无法监听端口: %v", err)
    }
    log.Println("gRPC 服务器启动成功,监听 :50051")
    // 启动服务器
    if err := server.Serve(lis); err != nil {
        log.Fatalf("服务器启动失败: %v", err)
    }
}

客户端实现 (client/main.go)

package main
import (
    "context"
    "io"
    "log"
    "time"
    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials/insecure"
    pb "grpc-demo/proto/user"
)
func main() {
    // 建立连接
    conn, err := grpc.Dial("localhost:50051", grpc.WithTransportCredentials(insecure.NewCredentials()))
    if err != nil {
        log.Fatalf("连接失败: %v", err)
    }
    defer conn.Close()
    // 创建客户端
    client := pb.NewUserServiceClient(conn)
    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()
    // 1. 一元 RPC 调用
    log.Println("=== 1. 获取单个用户 ===")
    getUser(client, ctx, 1)
    // 2. 服务端流式 RPC
    log.Println("\n=== 2. 获取用户列表(流式) ===")
    listUsers(client, ctx)
    // 3. 客户端流式 RPC
    log.Println("\n=== 3. 批量创建用户(客户端流式) ===")
    batchCreateUsers(client, ctx)
    // 4. 双向流式 RPC
    log.Println("\n=== 4. 用户聊天(双向流式) ===")
    chat(client, ctx)
}
// 一元 RPC 调用
func getUser(client pb.UserServiceClient, ctx context.Context, id int32) {
    req := &pb.GetUserRequest{Id: id}
    user, err := client.GetUser(ctx, req)
    if err != nil {
        log.Printf("获取用户失败: %v", err)
        return
    }
    log.Printf("获取用户成功: ID=%d, Name=%s, Email=%s, CreatedAt=%s",
        user.Id, user.Name, user.Email, user.CreatedAt)
}
// 服务端流式 RPC 调用
func listUsers(client pb.UserServiceClient, ctx context.Context) {
    req := &pb.ListUsersRequest{
        Page:     1,
        PageSize: 2,
    }
    stream, err := client.ListUsers(ctx, req)
    if err != nil {
        log.Printf("获取用户列表失败: %v", err)
        return
    }
    log.Println("正在接收用户列表...")
    for {
        user, err := stream.Recv()
        if err == io.EOF {
            log.Println("用户列表接收完成")
            break
        }
        if err != nil {
            log.Printf("接收用户数据失败: %v", err)
            break
        }
        log.Printf("收到用户: ID=%d, Name=%s, Email=%s",
            user.Id, user.Name, user.Email)
    }
}
// 客户端流式 RPC 调用
func batchCreateUsers(client pb.UserServiceClient, ctx context.Context) {
    stream, err := client.BatchCreateUsers(ctx)
    if err != nil {
        log.Printf("批量创建用户失败: %v", err)
        return
    }
    // 发送多个用户数据
    usersToCreate := []*pb.CreateUserRequest{
        {Name: "王五", Email: "wangwu@example.com"},
        {Name: "赵六", Email: "zhaoliu@example.com"},
        {Name: "孙七", Email: "sunqi@example.com"},
    }
    for _, req := range usersToCreate {
        log.Printf("发送创建请求: Name=%s", req.Name)
        if err := stream.Send(req); err != nil {
            log.Printf("发送失败: %v", err)
            return
        }
        time.Sleep(100 * time.Millisecond)
    }
    // 关闭流并获取响应
    resp, err := stream.CloseAndRecv()
    if err != nil {
        log.Printf("接收响应失败: %v", err)
        return
    }
    log.Printf("批量创建成功: 共创建 %d 个用户", resp.SuccessCount)
    for _, user := range resp.Users {
        log.Printf("  创建的用户: ID=%d, Name=%s, Email=%s", user.Id, user.Name, user.Email)
    }
}
// 双向流式 RPC 调用
func chat(client pb.UserServiceClient, ctx context.Context) {
    stream, err := client.Chat(ctx)
    if err != nil {
        log.Printf("聊天连接失败: %v", err)
        return
    }
    // 发送消息的 goroutine
    go func() {
        messages := []string{
            "你好,服务端!",
            "这是我的第二条消息",
            "你收到我的消息了吗?",
        }
        for _, msg := range messages {
            msg := &pb.ChatMessage{
                UserName:  "Client-User",
                Content:   msg,
                Timestamp: time.Now().Format("2006-01-02 15:04:05"),
            }
            if err := stream.Send(msg); err != nil {
                log.Printf("发送消息失败: %v", err)
                return
            }
            time.Sleep(500 * time.Millisecond)
        }
        // 关闭发送方向
        if err := stream.CloseSend(); err != nil {
            log.Printf("关闭发送失败: %v", err)
        }
    }()
    // 接收回复
    for {
        reply, err := stream.Recv()
        if err == io.EOF {
            log.Println("聊天结束")
            break
        }
        if err != nil {
            log.Printf("接收回复失败: %v", err)
            break
        }
        log.Printf("[%s] %s: %s", reply.Timestamp, reply.UserName, reply.Content)
    }
}

初始化 go.mod

module grpc-demo
go 1.19
require (
    google.golang.org/grpc v1.56.2
    google.golang.org/protobuf v1.31.0
)

运行测试

# 启动服务器
go run server/main.go
# 在另一个终端运行客户端
go run client/main.go

输出示例

=== 1. 获取单个用户 ===
获取用户成功: ID=1, Name=张三, Email=zhangsan@example.com, CreatedAt=2024-01-01 12:00:00
=== 2. 获取用户列表(流式) ===
正在接收用户列表...
收到用户: ID=1, Name=张三, Email=zhangsan@example.com
收到用户: ID=2, Name=李四, Email=lisi@example.com
用户列表接收完成
=== 3. 批量创建用户(客户端流式) ===
发送创建请求: Name=王五
发送创建请求: Name=赵六
发送创建请求: Name=孙七
批量创建成功: 共创建 3 个用户
  创建的用户: ID=3, Name=王五, Email=wangwu@example.com
  创建的用户: ID=4, Name=赵六, Email=zhaoliu@example.com
  创建的用户: ID=5, Name=孙七, Email=sunqi@example.com
=== 4. 用户聊天(双向流式) ===
[2024-01-01 12:00:01] Server: 收到消息: "你好,服务端!"
[2024-01-01 12:00:01] Server: 收到消息: "这是我的第二条消息"
[2024-01-01 12:00:02] Server: 收到消息: "你收到我的消息了吗?"
聊天结束

这个完整的案例展示了 gRPC 的四种通信模式:

  1. 一元 RPC:简单的请求-响应模式
  2. 服务端流式:服务端多次发送数据
  3. 客户端流式:客户端多次发送数据
  4. 双向流式:双方可以同时发送和接收数据

您可以根据实际需求修改和扩展这个案例。

抱歉,评论功能暂时关闭!