本文目录导读:

我将为您提供一个完整的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 的四种通信模式:
- 一元 RPC:简单的请求-响应模式
- 服务端流式:服务端多次发送数据
- 客户端流式:客户端多次发送数据
- 双向流式:双方可以同时发送和接收数据
您可以根据实际需求修改和扩展这个案例。