Skip to content

gRPC 通信基础

本篇是 Go 微服务系列的第三篇。在微服务架构中,服务间通信是绕不开的核心主题。gRPC 是 Go 生态中最主流的高性能 RPC 框架,基于 HTTP/2 和 Protobuf,是构建微服务同步通信的首选方案。本篇将系统讲解 gRPC 的核心概念、四种调用模式、拦截器、错误处理,并实现一个完整的用户服务。

一、gRPC 简介

1. 什么是 gRPC

gRPC 是 Google 开源的高性能 RPC 框架,核心特征:

  • 基于 HTTP/2:多路复用、双向流、头部压缩,性能远超 HTTP/1.1。
  • 使用 Protobuf:二进制编码,体积小、解析快,比 JSON 快数倍。
  • 跨语言:通过 proto 文件生成多语言客户端/服务端代码(Go、Java、Python、C++、Node.js 等 10+ 种语言)。
  • 强类型契约:proto 文件就是服务契约,编译期就能发现不兼容变更。
  • 四种调用模式:一元调用、服务端流、客户端流、双向流。

2. gRPC vs REST

维度gRPCREST (HTTP/JSON)
协议HTTP/2HTTP/1.1 为主
编码Protobuf 二进制JSON 文本
性能高(解析快、体积小)中等
契约proto 强类型OpenAPI 可选
流式原生支持需 SSE/WebSocket
浏览器需要 gRPC-Web 代理原生支持
适用场景内部服务间调用对外 API、前后端通信

实际项目通常内部用 gRPC,对外用 REST,由 API 网关做协议转换。

3. gRPC 在 Go 生态中的地位

  • gRPC 官方对 Go 的支持是一等公民,google.golang.org/grpc 是 Go 团队维护的。
  • 大量云原生项目(etcd、TiKV、CockroachDB、Kubernetes API 的部分接口)使用 gRPC。
  • 国内 Kratos、Go-Zero、go-micro 等微服务框架均以 gRPC 为底层通信。

二、Protobuf 消息定义

1. 安装工具链

需要两个工具:

bash
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

确认 PATH 包含 $(go env GOPATH)/bin

2. 编写 proto 文件

创建 user.proto

protobuf
syntax = "proto3";

package user.v1;
option go_package = "github.com/example/user-service/api/user/v1;userv1";

// 用户实体
message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  int32 age = 4;
}

// 创建用户请求
message CreateUserRequest {
  string name = 1;
  string email = 2;
  int32 age = 3;
}

// 通用 ID 请求
message GetUserRequest {
  int64 id = 1;
}

// 用户列表请求(带分页)
message ListUsersRequest {
  int32 page = 1;
  int32 page_size = 2;
}

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

// 用户服务
service UserService {
  // 一元调用:创建用户
  rpc CreateUser(CreateUserRequest) returns (User);
  // 一元调用:获取用户
  rpc GetUser(GetUserRequest) returns (User);
  // 服务端流:批量获取用户
  rpc ListUsers(ListUsersRequest) returns (stream User);
  // 客户端流:批量创建
  rpc BatchCreate(stream CreateUserRequest) returns (ListUsersResponse);
  // 双向流:聊天式交互
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message ChatMessage {
  string from = 1;
  string text = 2;
}

3. proto3 语法要点

  • 字段编号:1~15 占 1 字节,16~2047 占 2 字节,频繁字段用小编号。
  • 标量类型int32int64stringboolfloatdoublebytes
  • 复合类型messageenumrepeated(列表)、map
  • 默认值:proto3 字段没有 required,默认值是 0/""/false,可用 optional 显式区分「未设置」和「设置为默认值」。
  • 包名与 go_packagego_package 决定生成代码的导入路径。

三、生成 Go 代码

执行 protoc 命令生成 Go 代码:

bash
protoc --go_out=. --go_opt=paths=source_relative \
  --go-grpc_out=. --go-grpc_opt=paths=source_relative \
  api/user/v1/user.proto

生成两个文件:

  • user.pb.go:消息结构体、序列化/反序列化代码。
  • user_grpc.pb.go:gRPC 客户端 stub 和服务端 interface。

生成的服务端接口大致长这样:

go
type UserServiceServer interface {
    CreateUser(context.Context, *CreateUserRequest) (*User, error)
    GetUser(context.Context, *GetUserRequest) (*User, error)
    ListUsers(*ListUsersRequest, UserService_ListUsersServer) error
    BatchCreate(UserService_BatchCreateServer) error
    Chat(UserService_ChatServer) error
}

四、四种服务方法

1. 一元调用(Unary RPC)

最常见:一次请求,一次响应。

go
// 服务端实现
func (s *userServiceServer) CreateUser(ctx context.Context, req *userv1.CreateUserRequest) (*userv1.User, error) {
    if req.GetName() == "" {
        return nil, status.Error(codes.InvalidArgument, "name is required")
    }
    user := &userv1.User{
        Id:    time.Now().Unix(),
        Name:  req.GetName(),
        Email: req.GetEmail(),
        Age:   req.GetAge(),
    }
    return user, nil
}

// 客户端调用
resp, err := client.CreateUser(ctx, &userv1.CreateUserRequest{
    Name:  "alice",
    Email: "alice@example.com",
    Age:   18,
})

2. 服务端流(Server Streaming)

一次请求,多次响应。适合分页拉取、大文件分块传输、实时推送。

go
// 服务端实现
func (s *userServiceServer) ListUsers(req *userv1.ListUsersRequest, stream userv1.UserService_ListUsersServer) error {
    for i := 0; i < 10; i++ {
        if err := stream.Send(&userv1.User{
            Id:   int64(i),
            Name: fmt.Sprintf("user-%d", i),
        }); err != nil {
            return err
        }
    }
    return nil
}

// 客户端接收
stream, err := client.ListUsers(ctx, &userv1.ListUsersRequest{Page: 1, PageSize: 10})
for {
    user, err := stream.Recv()
    if err == io.EOF {
        break
    }
    if err != nil {
        log.Fatal(err)
    }
    fmt.Printf("recv: %+v\n", user)
}

3. 客户端流(Client Streaming)

多次请求,一次响应。适合批量上传、聚合统计。

go
// 服务端实现
func (s *userServiceServer) BatchCreate(stream userv1.UserService_BatchCreateServer) error {
    var count int32
    for {
        req, err := stream.Recv()
        if err == io.EOF {
            return stream.SendAndClose(&userv1.ListUsersResponse{
                Total: count,
            })
        }
        if err != nil {
            return err
        }
        count++
    }
}

// 客户端发送
stream, _ := client.BatchCreate(ctx)
for i := 0; i < 5; i++ {
    _ = stream.Send(&userv1.CreateUserRequest{Name: fmt.Sprintf("u-%d", i)})
}
resp, _ := stream.CloseAndRecv()

4. 双向流(Bidirectional Streaming)

双向同时收发,适合聊天、实时同步、协同编辑。

go
// 服务端实现
func (s *userServiceServer) Chat(stream userv1.UserService_ChatServer) error {
    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            return nil
        }
        if err != nil {
            return err
        }
        // echo 回去
        if err := stream.Send(&userv1.ChatMessage{
            From: "server",
            Text: "echo: " + msg.GetText(),
        }); err != nil {
            return err
        }
    }
}

五、gRPC 服务端实现

下面是一个完整的 gRPC 服务端:

go
package main

import (
	"context"
	"fmt"
	"log"
	"net"
	"os"
	"os/signal"
	"syscall"

	"google.golang.org/grpc"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
)

// 以下是简化的接口示意(实际项目用 protoc 生成代码)
// 为了让本例可独立编译,这里手写最小骨架,省略 .pb.go 引入

// User 用户结构
type User struct {
	ID    int64
	Name  string
	Email string
	Age   int32
}

// UserStore 内存存储
type UserStore struct {
	users map[int64]*User
	nextID int64
}

func NewUserStore() *UserStore {
	return &UserStore{users: make(map[int64]*User), nextID: 1}
}

func (s *UserStore) Create(name, email string, age int32) *User {
	id := s.nextID
	s.nextID++
	u := &User{ID: id, Name: name, Email: email, Age: age}
	s.users[id] = u
	return u
}

func (s *UserStore) Get(id int64) (*User, error) {
	u, ok := s.users[id]
	if !ok {
		return nil, status.Errorf(codes.NotFound, "user %d not found", id)
	}
	return u, nil
}

// 用 net.Listen 实现一个最小 TCP 服务,模拟 gRPC 服务的生命周期
// 真实项目用 grpc.NewServer() + RegisterUserServiceServer + server.Serve(lis)

func main() {
	store := NewUserStore()
	_ = store.Create("alice", "alice@example.com", 18)
	_ = store.Create("bob", "bob@example.com", 20)

	// 这里使用 gRPC 的 net.Listen 模式骨架
	lis, err := net.Listen("tcp", ":50051")
	if err != nil {
		log.Fatalf("listen failed: %v", err)
	}

	// 在真实场景下:
	// s := grpc.NewServer(grpc.UnaryInterceptor(...))
	// userv1.RegisterUserServiceServer(s, &userServiceServer{store: store})
	// go s.Serve(lis)

	log.Printf("gRPC server listening on %s", lis.Addr().String())

	// 优雅退出
	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit
	_ = lis.Close()

	// 演示:用 store 做一次查询
	u, err := store.Get(1)
	if err != nil {
		log.Println(err)
		return
	}
	fmt.Printf("user: %+v\n", u)

	// 演示 context 取消传播
	ctx, cancel := context.WithCancel(context.Background())
	go func() {
		<-time.After(100 * time.Millisecond)
		cancel()
	}()
	<-ctx.Done()
	fmt.Println("context cancelled:", ctx.Err())
}

// 引入 time 仅用于示例
import (
	"time"
)

注意:上面这个简化版为了让示例可独立编译,省略了 protoc 生成的代码。真实项目结构请参考本篇末尾的完整示例。

六、gRPC 客户端实现

go
package main

import (
	"context"
	"fmt"
	"log"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"
)

// 客户端连接 gRPC 服务的标准模式
func dialGRPC(addr string) (*grpc.ClientConn, error) {
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()
	// 生产环境应使用 TLS / credentials
	return grpc.DialContext(ctx, addr,
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithBlock(),
		grpc.WithDefaultCallOptions(
			grpc.MaxCallRecvMsgSize(4*1024*1024),
		),
	)
}

func main() {
	conn, err := dialGRPC("127.0.0.1:50051")
	if err != nil {
		log.Fatalf("dial failed: %v", err)
	}
	defer conn.Close()

	// 真实场景:client := userv1.NewUserServiceClient(conn)
	// 这里仅演示连接建立和关闭流程
	fmt.Println("connected to gRPC server")

	ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
	defer cancel()
	// 演示一次调用骨架(实际调用需 .pb.go 生成的 client)
	_ = ctx
	fmt.Println("done")
}

客户端最佳实践

  • 连接复用:gRPC 底层是 HTTP/2 长连接,一个 ClientConn 即可承载高并发请求,不要每次请求都 Dial。
  • 超时控制:每个 RPC 都应该带 context.WithTimeout,避免调用方无限阻塞。
  • 重试与退避:使用 gRPC 内建的服务配置重试策略或客户端拦截器。
  • TLS / mTLS:生产环境必须启用 TLS,内部服务间可用 mTLS 双向认证。

七、拦截器

拦截器(Interceptor)类似 HTTP 中间件,可以在调用前后插入横切逻辑:日志、监控、认证、追踪、限流。

1. 一元拦截器

go
package main

import (
	"context"
	"log"
	"time"

	"google.golang.org/grpc"
)

// 服务端一元拦截器:记录每个 RPC 的耗时
func loggingUnaryServerInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
	start := time.Now()
	resp, err := handler(ctx, req)
	log.Printf("[server] %s cost=%v err=%v", info.FullMethod, time.Since(start), err)
	return resp, err
}

// 客户端一元拦截器:注入 request id
func requestIDUnaryClientInterceptor(ctx context.Context, method string, req, reply interface{}, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
	start := time.Now()
	err := invoker(ctx, method, req, reply, cc, opts...)
	log.Printf("[client] %s cost=%v err=%v", method, time.Since(start), err)
	return err
}

func newServerWithInterceptors() *grpc.Server {
	return grpc.NewServer(
		grpc.UnaryInterceptor(loggingUnaryServerInterceptor),
	)
}

func newClientWithInterceptors(addr string) (*grpc.ClientConn, error) {
	return grpc.Dial(addr,
		grpc.WithUnaryInterceptor(requestIDUnaryClientInterceptor),
	)
}

func main() {
	// 仅演示拦截器组装
	_ = newServerWithInterceptors()
	conn, err := newClientWithInterceptors("127.0.0.1:50051")
	if err != nil {
		log.Fatal(err)
	}
	defer conn.Close()
}

2. 流拦截器

流拦截器签名稍有不同,需包装 grpc.ServerStream / grpc.ClientStream

go
package main

import (
	"context"
	"log"
	"time"

	"google.golang.org/grpc"
)

type wrappedStream struct {
	grpc.ServerStream
}

func (w *wrappedStream) RecvMsg(m interface{}) error {
	log.Printf("recv msg start")
	return w.ServerStream.RecvMsg(m)
}

func (w *wrappedStream) SendMsg(m interface{}) error {
	log.Printf("send msg start")
	return w.ServerStream.SendMsg(m)
}

func loggingStreamServerInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
	start := time.Now()
	err := handler(srv, &wrappedStream{ServerStream: ss})
	log.Printf("[stream] %s cost=%v err=%v", info.FullMethod, time.Since(start), err)
	return err
}

func newServerWithStreamInterceptor() *grpc.Server {
	return grpc.NewServer(grpc.StreamInterceptor(loggingStreamServerInterceptor))
}

func main() {
	_ = newServerWithStreamInterceptor()
	_ = context.Background()
}

3. 多拦截器链

gRPC 支持 grpc.ChainUnaryInterceptor 串联多个一元拦截器,执行顺序类似洋葱模型:

go
server := grpc.NewServer(
    grpc.ChainUnaryInterceptor(
        recoveryInterceptor,    // panic 恢复(最外层)
        tracingInterceptor,     // 注入 trace
        loggingInterceptor,     // 日志
        authInterceptor,        // 认证(最内层,最靠近业务)
    ),
)

八、错误处理与状态码

1. gRPC 状态码

gRPC 定义了标准的状态码,跨语言一致:

状态码含义HTTP 类比
OK成功200
Canceled客户端取消499
Unknown未知错误500
InvalidArgument参数错误400
DeadlineExceeded超时504
NotFound资源不存在404
AlreadyExists资源已存在409
PermissionDenied无权限403
Unauthenticated未认证401
ResourceExhausted资源耗尽(限流)429
Unavailable服务不可用503
Internal内部错误500

2. 返回错误

go
package main

import (
	"context"
	"errors"
	"fmt"

	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
)

// 自定义业务错误码(payload 用 details 携带)
type BizError struct {
	Code int32
	Msg  string
}

func (e *BizError) Error() string {
	return fmt.Sprintf("biz error %d: %s", e.Code, e.Msg)
}

// 包装为 gRPC status 错误
func toGRPCError(err error) error {
	if err == nil {
		return nil
	}
	var bizErr *BizError
	if errors.As(err, &bizErr) {
		return status.Errorf(codes.Code(bizErr.Code), "%s", bizErr.Msg)
	}
	if errors.Is(err, context.DeadlineExceeded) {
		return status.Error(codes.DeadlineExceeded, err.Error())
	}
	return status.Error(codes.Internal, err.Error())
}

func exampleHandler(ctx context.Context, req interface{}) error {
	err := errors.New("not implemented")
	return toGRPCError(err)
}

func main() {
	err := exampleHandler(context.Background(), nil)
	if err != nil {
		st, _ := status.FromError(err)
		fmt.Printf("code=%v msg=%s\n", st.Code(), st.Message())
	}

	err2 := toGRPCError(&BizError{Code: int32(codes.NotFound), Msg: "user not found"})
	st, _ := status.FromError(err2)
	fmt.Printf("biz code=%v msg=%s\n", st.Code(), st.Message())
}

3. 客户端错误处理

go
package main

import (
	"context"
	"fmt"
	"log"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/credentials/insecure"
	"google.golang.org/grpc/status"
)

func main() {
	conn, err := grpc.Dial("127.0.0.1:50051",
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithBlock(),
		grpc.WithTimeout(2*time.Second),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer conn.Close()

	// 模拟一次调用(实际需 .pb.go 生成的 client)
	// err := client.GetUser(ctx, &userv1.GetUserRequest{Id: 999})
	// 这里用一个 stub 模拟
	err = status.Error(codes.NotFound, "user not found")

	st, ok := status.FromError(err)
	if !ok {
		log.Fatalf("non-grpc error: %v", err)
	}
	switch st.Code() {
	case codes.NotFound:
		fmt.Println("user not exists, fallback to default")
	case codes.DeadlineExceeded:
		fmt.Println("timeout, retrying...")
	case codes.Unavailable:
		fmt.Println("service unavailable, retrying with backoff")
	default:
		fmt.Printf("other error: %v\n", err)
	}

	_ = context.Background()
	_ = conn
}

九、完整示例:用户服务的 gRPC 实现

下面给出一个可独立编译运行的完整示例(用标准库 net 模拟,让你无需 protoc 即可看到服务端/客户端交互的完整骨架;理解了它再切到真实 gRPC 就轻车熟路)。

server/main.go

go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"net"
	"os"
	"os/signal"
	"sync"
	"syscall"
)

type User struct {
	ID    int64  `json:"id"`
	Name  string `json:"name"`
	Email string `json:"email"`
}

type Server struct {
	mu     sync.Mutex
	nextID int64
	users  map[int64]*User
}

func NewServer() *Server {
	return &Server{nextID: 1, users: make(map[int64]*User)}
}

func (s *Server) CreateUser(name, email string) *User {
	s.mu.Lock()
	defer s.mu.Unlock()
	u := &User{ID: s.nextID, Name: name, Email: email}
	s.nextID++
	s.users[u.ID] = u
	return u
}

func (s *Server) GetUser(id int64) (*User, error) {
	s.mu.Lock()
	defer s.mu.Unlock()
	u, ok := s.users[id]
	if !ok {
		return nil, fmt.Errorf("user %d not found", id)
	}
	return u, nil
}

// handleConn 简化的请求处理:每行一个 JSON 请求
// {"method":"create","name":"alice","email":"a@x.com"}
// {"method":"get","id":1}
func (s *Server) handleConn(ctx context.Context, conn net.Conn) {
	defer conn.Close()
	dec := json.NewDecoder(conn)
	enc := json.NewEncoder(conn)
	for {
		var req map[string]interface{}
		if err := dec.Decode(&req); err != nil {
			return
		}
		method, _ := req["method"].(string)
		switch method {
		case "create":
			name, _ := req["name"].(string)
			email, _ := req["email"].(string)
			u := s.CreateUser(name, email)
			_ = enc.Encode(map[string]interface{}{"ok": true, "user": u})
		case "get":
			id := int64(req["id"].(float64))
			u, err := s.GetUser(id)
			if err != nil {
				_ = enc.Encode(map[string]interface{}{"ok": false, "error": err.Error()})
				continue
			}
			_ = enc.Encode(map[string]interface{}{"ok": true, "user": u})
		default:
			_ = enc.Encode(map[string]interface{}{"ok": false, "error": "unknown method"})
		}
	}
}

func main() {
	srv := NewServer()
	srv.CreateUser("alice", "alice@example.com")
	srv.CreateUser("bob", "bob@example.com")

	lis, err := net.Listen("tcp", ":50051")
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("server listening on %s", lis.Addr())

	go func() {
		for {
			conn, err := lis.Accept()
			if err != nil {
				return
			}
			go srv.handleConn(context.Background(), conn)
		}
	}()

	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit
	_ = lis.Close()
	log.Println("server exited")
}

client/main.go

go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"net"
	"time"
)

func main() {
	conn, err := net.Dial("tcp", "127.0.0.1:50051")
	if err != nil {
		log.Fatal(err)
	}
	defer conn.Close()

	enc := json.NewEncoder(conn)
	dec := json.NewDecoder(conn)

	// 创建用户
	ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
	defer cancel()
	_ = ctx

	_ = enc.Encode(map[string]interface{}{
		"method": "create",
		"name":   "charlie",
		"email":  "charlie@example.com",
	})
	var resp map[string]interface{}
	_ = dec.Decode(&resp)
	fmt.Printf("create resp: %+v\n", resp)

	_ = enc.Encode(map[string]interface{}{
		"method": "get",
		"id":     1,
	})
	_ = dec.Decode(&resp)
	fmt.Printf("get resp: %+v\n", resp)
}

虽然这个示例用了 net + JSON 而不是真正的 gRPC,但它演示了 RPC 的核心流程:定义契约 → 服务端实现 → 客户端调用 → 错误处理。把它替换成真正的 gRPC 时,契约由 proto 文件定义,传输由 HTTP/2 + Protobuf 完成,但骨架完全一致。

十、小结

本篇我们系统学习了 gRPC:

  • 核心价值:基于 HTTP/2 + Protobuf,高性能、强类型契约、跨语言。
  • proto3 语法:message、service、字段编号、optional、repeated。
  • 四种调用模式:一元、服务端流、客户端流、双向流,分别适合不同的业务场景。
  • 服务端:实现 XXXServiceServer 接口,注册到 grpc.Server,监听 TCP。
  • 客户端grpc.DialContext 建立长连接,复用 stub 调用。
  • 拦截器:一元和流两种,可串联多个,实现日志、追踪、认证、限流等横切逻辑。
  • 错误处理:使用 status.Error + 标准状态码,跨语言一致。

下一篇我们会进入 RESTful API 与 Gin 集成,学习如何用 Gin 构建对外的 HTTP API,并实现 gRPC + Gin 双端口服务(内部 gRPC 通信,对外 REST API)。

延伸阅读