vault backup: 2026-05-11 19:02:38
This commit is contained in:
@@ -0,0 +1,426 @@
|
||||
---
|
||||
tags: [gRPC, Go, Server, Production]
|
||||
create time: 2026-05-11 16:00
|
||||
---
|
||||
|
||||
# Server 搭建与注册
|
||||
|
||||
## 概述
|
||||
|
||||
gRPC server 搭建本身很简单——`grpc.NewServer()` 加一行 `RegisterxxxServer()` 就够了。但在生产环境中你需要考虑的东西很多:TLS 配置、优雅关闭、服务发现集成、多端口暴露、reflection 开关、健康检查。我们从一个最简单的 hello world 开始,逐步构建一个 production-ready 的 server。
|
||||
|
||||
> [!question] 为什么生产环境需要关心 graceful shutdown?
|
||||
> gRPC 基于 HTTP/2,连接是长连接。如果直接 kill 进程,所有正在处理的请求会突然断掉,client 端收到的是 TCP RST 而非一个干净的 finish。这在金融或订单系统中可能导致重复扣款、状态不一致。
|
||||
|
||||
## 最小可运行 Server
|
||||
|
||||
```go
|
||||
// cmd/server/main.go
|
||||
package main
|
||||
|
||||
import (
|
||||
"log"
|
||||
"net"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/reflection"
|
||||
|
||||
pb "your/proto/gen/go"
|
||||
)
|
||||
|
||||
func main() {
|
||||
logged, _ := zap.NewProduction()
|
||||
zap.ReplaceGlobals(logged)
|
||||
defer logged.Sync()
|
||||
|
||||
lis, err := net.Listen("tcp", ":50051")
|
||||
if err != nil {
|
||||
log.Fatalf("failed to listen: %v", err)
|
||||
}
|
||||
|
||||
s := grpc.NewServer()
|
||||
pb.RegisterUserServiceServer(s, &userService{})
|
||||
|
||||
// Reflection for debugging — production 应关闭
|
||||
reflection.Register(s)
|
||||
|
||||
zap.L().Sugar().Infow("serving", "addr", lis.Addr())
|
||||
|
||||
if err := s.Serve(lis); err != nil {
|
||||
log.Fatalf("failed to serve: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// userService 实现 pb.UserServiceServer interface
|
||||
type userService struct {
|
||||
pb.UnimplementedUserServiceServer // 嵌入 zero-stub,天然支持未来 method 新增
|
||||
}
|
||||
```
|
||||
|
||||
核心就三步:监听端口 → 创建 server → 启动。**但** `UnimplementedUserServiceServer` 这个嵌入类型很重要——它让你无需为每个 method 写空壳,proto 文件新增 RPC method 时编译期自动报错提醒,而不是静默忽略。
|
||||
|
||||
## Server Options(ServerOption 精选)
|
||||
|
||||
`grpc.NewServer(...)` 接受任意数量的 `ServerOption`(函数式选项模式)。以下是高频选项速查表:
|
||||
|
||||
| Option | 用途 | 默认值 | 建议 |
|
||||
|--------|------|--------|------|
|
||||
| `grpc.Creds(credentials)` | TLS / mTLS 认证 | 明文 | **生产必配** |
|
||||
| `grpc.MaxRecvMsgSize(n int)` | 单条接收消息上限 | 4MB | 按需调大,注意内存 |
|
||||
| `grpc.MaxSendMsgSize(n int)` | 单条发送消息上限 | 4MB | 配合 MaxRecv 对称设置 |
|
||||
| `grpc.MaxConcurrentStreams(n uint32)` | 单 stream 最大并发 HEADERS | 100 | 防 DoS,放宽到 1024+ |
|
||||
| `grpc.KeepaliveParams(kp)` | keepalive 心跳参数 | 2h idle timeout | LB 后必须调整 |
|
||||
| `grpc.ChainUnaryInterceptor(ints...)` | Unary 拦截器链 | 无 | 日志 / 鉴权 / 追踪 |
|
||||
| `grpc.StatsHandler(h stats.Handler)` | OpenTelemetry 等 stats 注入 | 无 | 链路追踪 / metrics |
|
||||
|
||||
```go
|
||||
s := grpc.NewServer(
|
||||
grpc.Creds(credentials.NewTLS(tlsConfig)),
|
||||
grpc.MaxRecvMsgSize(16*1024*1024), // 16MB
|
||||
grpc.MaxConcurrentStreams(1024), // 放宽并发流限制
|
||||
grpc.KeepaliveParams(keepalive.ServerParameters{
|
||||
MaxConnectionIdle: 15 * time.Minute,
|
||||
KeepaliveTime: 20 * time.Second,
|
||||
KeepaliveTimeout: 5 * time.Second,
|
||||
MinTimeBetweenPings: 10 * time.Second,
|
||||
PingWithoutCallsAllowed: true,
|
||||
}),
|
||||
grpc.ChainUnaryInterceptor(logging.Unary(), auth.Unary()),
|
||||
// grpc.StatsHandler(otelgrpc.NewServerHandler()), // OpenTelemetry Go SDK
|
||||
)
|
||||
```
|
||||
|
||||
> [!tip] keepalive 参数调优经验
|
||||
> 经过 LB(如 Envoy 或 Nginx)时,默认的 2 小时 idle timeout 会让 LB 提前切断连接,导致 client 侧出现 "connection reset" 错误。把 `MaxConnectionIdle` 设短到 15~20 分钟更合理——让 server 主动重建连接,确保两端对连接生命周期有一致认知。
|
||||
>
|
||||
> 直连无 LB 场景可以保持默认值,减少不必要的连接重建开销。
|
||||
|
||||
### 拦截器链执行顺序
|
||||
|
||||
拦截器以**尾递归**方式组合——先注册的 interceptor 包裹后注册的,形成洋葱模型:
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
Client["Client Request"] --> Logging["Logging Interceptor<br/>(最外层)"]
|
||||
Logging --> Auth["Auth Interceptor<br/>(内层)"]
|
||||
Auth --> Recovery["Recovery Interceptor<br/>(最内层)"]
|
||||
Recovery --> Handler["Actual Handler"]
|
||||
|
||||
style Client fill:#E3F2FD
|
||||
style Handler fill:#FFF3E0
|
||||
```
|
||||
|
||||
```go
|
||||
// 请求到达顺序:Logging → Auth → Recovery → Handler
|
||||
// 响应返回顺序:Handler → Recovery → Auth → Logging
|
||||
grpc.ChainUnaryInterceptor(
|
||||
logging.UnaryInterceptor(), // 第 1 层(最外)
|
||||
auth.UnaryInterceptor(), // 第 2 层
|
||||
recovery.UnaryInterceptor(), // 第 3 层(最内)
|
||||
)
|
||||
```
|
||||
|
||||
**设计建议:**
|
||||
- **外层**做横切关注点:日志、追踪、metrics
|
||||
- **中层**做业务安全校验:鉴权、限流
|
||||
- **内层**做兜底逻辑:panic recover、超时控制
|
||||
|
||||
## Service Registration
|
||||
|
||||
### 什么是 Service Registration
|
||||
|
||||
每次 proto 文件中的 `service` 块都会生成两个 Go 类型:
|
||||
|
||||
1. **`<ServiceName>Server` interface** —— 定义了你必须实现的 RPC method 集合
|
||||
2. **`Register<ServiceName>Server(server, impl)` 函数** —— 将实现注册到 server 的 method dispatch table
|
||||
|
||||
```go
|
||||
// 你定义的 .proto:
|
||||
// service UserService {
|
||||
// rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
|
||||
// rpc GetUser(GetUserRequest) returns (GetUserResponse);
|
||||
// }
|
||||
|
||||
// 生成的接口(简化示意):
|
||||
type UserServiceServer interface {
|
||||
CreateUser(context.Context, *pb.CreateUserRequest) (*pb.CreateUserResponse, error)
|
||||
GetUser(context.Context, *pb.GetUserRequest) (*pb.GetUserResponse, error)
|
||||
// UnimplementedUserServiceServer 提供 default implementation (Unimplemented)
|
||||
mustEmbedUnimplementedUserServiceServer()
|
||||
}
|
||||
|
||||
// 注册函数签名:
|
||||
func RegisterUserServiceServer(s *grpc.Server, srv UserServiceServer)
|
||||
```
|
||||
|
||||
### 单 server 注册多 service
|
||||
|
||||
不需要为每个 service 创建独立 `grpc.Server`。单个 `grpc.Server` 实例天然支持多个 service registration:
|
||||
|
||||
```go
|
||||
s := grpc.NewServer(opts...)
|
||||
|
||||
pb.RegisterUserServiceServer(s, NewUserService())
|
||||
pb.RegisterOrderServiceServer(s, NewOrderService())
|
||||
pb.RegisterPaymentServiceServer(s, NewPaymentService())
|
||||
|
||||
s.Serve(lis) // 一个 listener 对外暴露三个 service
|
||||
```
|
||||
|
||||
底层通过 method name 路由(例如 `:method/user.v1.UserService/CreateUser`),各 service 完全隔离互不干扰。
|
||||
|
||||
### reflection:开发调试 vs 生产关闭
|
||||
|
||||
Reflection 允许外部工具在运行时查询已注册的 service descriptor,无需 `.proto` 文件:
|
||||
|
||||
```go
|
||||
import "google.golang.org/grpc/reflection"
|
||||
|
||||
reflection.Register(s) // 注册 gRPC Reflection v1alpha 服务
|
||||
```
|
||||
|
||||
**开发阶段好处:**
|
||||
- `grpcurl` 无需 `.proto` 即可列出可用 API
|
||||
- IDE 自动补全 gRPC call
|
||||
- 快速验证某个 service 是否成功注册
|
||||
|
||||
**生产环境关闭原因:**
|
||||
- 暴露了完整的 API schema 信息(方法名、消息结构),增加攻击面
|
||||
- 额外占用内存保存所有 descriptor protobuf
|
||||
- 客户端可通过 reflection 探测内部方法名(即使是未公开的)
|
||||
|
||||
```go
|
||||
if os.Getenv("ENV") != "production" {
|
||||
reflection.Register(s) // 仅非生产环境开启
|
||||
}
|
||||
```
|
||||
|
||||
## Multi-Port / Multi-Network
|
||||
|
||||
常见需求:内网接口和外网接口分开监听,或者同时暴露 IPv4 和 IPv6:
|
||||
|
||||
```go
|
||||
srv := grpc.NewServer(opts...)
|
||||
pb.RegisterUserServiceServer(srv, svc)
|
||||
|
||||
lis1, _ := net.Listen("tcp", ":50051") // 内网
|
||||
lis2, _ := net.Listen("tcp", ":50052") // 外网
|
||||
|
||||
go func() { _ = srv.Serve(lis1) }()
|
||||
go func() { _ = srv.Serve(lis2) }()
|
||||
|
||||
<-ctx.Done()
|
||||
```
|
||||
|
||||
> [!warning] 同一个 `grpc.Server` 不能在同一时刻被多个 goroutine 同时调用 `Serve()`(不安全)
|
||||
> 上面的写法在 gRPC-Go 中实际会导致 race condition。如果你的业务是按服务拆分到不同端口的需求,应该创建独立的 server 实例:
|
||||
|
||||
```go
|
||||
// 正确做法:每个端口使用独立的 grpc.Server 实例
|
||||
srv1 := grpc.NewServer(opts...)
|
||||
pb.RegisterUserServiceServer(srv1, userSvc)
|
||||
|
||||
srv2 := grpc.NewServer(opts...)
|
||||
pb.RegisterUserServiceServer(srv2, userSvc)
|
||||
pb.RegisterOrderServiceServer(srv2, orderSvc) // 外网额外暴露 OrderService
|
||||
|
||||
go func() { _ = srv1.Serve(net.Listen("tcp", ":50051")) }() // 内网:只暴露 User
|
||||
go func() { _ = srv2.Serve(net.Listen("tcp", ":50052")) }() // 外网:暴露 User + Order
|
||||
|
||||
<-ctx.Done()
|
||||
srv1.Stop()
|
||||
srv2.Stop()
|
||||
```
|
||||
|
||||
这样每个 server 实例管理自己的 listener 和连接池,彼此隔离互不影响。按服务粒度差异化暴露端口是微服务架构中的常见模式——内网服务只暴露核心 RPC,网关层再聚合多服务。
|
||||
|
||||
## Graceful Shutdown(重要!)
|
||||
|
||||
这是最容易踩坑的部分。gRPC 提供了两种停止方法:
|
||||
|
||||
```go
|
||||
// 优雅停止:拒绝新连接 + 等待活跃流完成
|
||||
s.GracefulStop()
|
||||
|
||||
// 立即停止:直接断开所有连接(未完成请求会报 error)
|
||||
s.Stop()
|
||||
```
|
||||
|
||||
典型的生产模式是结合 signal handling + context 超时控制:
|
||||
|
||||
```go
|
||||
func run(ctx context.Context, s *grpc.Server) error {
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
log.Println("shutdown signal received")
|
||||
s.GracefulStop()
|
||||
}()
|
||||
|
||||
return s.Serve(lis)
|
||||
}
|
||||
|
||||
// 30s 超时: GracefulStop 后最多等 30s,超时则强制退出
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := run(ctx, s); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
```
|
||||
|
||||
**为什么需要 30s 超时?** `GracefulStop` 没有内置超时机制——如果某个 handler goroutine 因为缺少 `ctx.Done()` 检测而永远不返回,它会阻塞整个关闭流程。超时兜底确保进程最终能退出(操作系统会重新发送 SIGKILL)。
|
||||
|
||||
GracefulStop 的执行流程如下:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
SIGTERM["SIGTERM Signal"] --> Handler["Signal Handler"]
|
||||
Handler --> StopNew["停止接收新连接<br/>新请求返回 UNAVAILABLE"]
|
||||
StopNew --> Drain["排空活跃 Stream<br/>等待 handler 返回"]
|
||||
Drain --> Check{"所有流<br/>全部完成?"}
|
||||
Check -->|"超时"| Force["强制关闭<br/>可能丢失未完成请求"]
|
||||
Check -->|"是"| CleanExit["干净退出<br/>所有回调执行完毕"]
|
||||
|
||||
style CleanExit fill:#00D866,color:#fff
|
||||
style Force fill:#EE5A24,color:#fff
|
||||
```
|
||||
|
||||
> [!danger] GracefulStop 不是万能的
|
||||
> 如果你的 handler 中没有正确检测 `ctx.Done()`,handler goroutine 永远不会结束,`GracefulStop` 最终也会卡住。所以 streaming handler 里必须遵循 ctx 驱动模式——见下一篇文章。
|
||||
|
||||
## 完整生产级模板
|
||||
|
||||
汇总成一个可直接复用的 bootstrap 骨架:
|
||||
|
||||
```go
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"google.golang.org/grpc"
|
||||
grpc_health_v1 "google.golang.org/grpc/health/grpc_health_v1"
|
||||
"google.golang.org/grpc/health"
|
||||
"google.golang.org/grpc/keepalive"
|
||||
"google.golang.org/grpc/reflection"
|
||||
|
||||
pb "your/proto/gen/go"
|
||||
)
|
||||
|
||||
func NewGRPCServer(opts ...grpc.ServerOption) (*grpc.Server, error) {
|
||||
// 初始化 logger(defer Sync 由调用方管理生命周期)
|
||||
logger, _ := zap.NewProduction()
|
||||
zap.ReplaceGlobals(logger)
|
||||
|
||||
creds, err := loadTLS()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("tls: %w", err)
|
||||
}
|
||||
|
||||
allOpts := []grpc.ServerOption{
|
||||
grpc.Creds(creds),
|
||||
grpc.MaxRecvMsgSize(16 * 1024 * 1024),
|
||||
grpc.MaxConcurrentStreams(1024),
|
||||
grpc.KeepaliveParams(keepalive.ServerParameters{
|
||||
MaxConnectionIdle: 15 * time.Minute,
|
||||
KeepaliveTime: 20 * time.Second,
|
||||
KeepaliveTimeout: 5 * time.Second,
|
||||
}),
|
||||
grpc.ChainUnaryInterceptor(
|
||||
logging.UnaryInterceptor(),
|
||||
auth.UnaryInterceptor(),
|
||||
),
|
||||
}
|
||||
allOpts = append(allOpts, opts...)
|
||||
|
||||
srv := grpc.NewServer(allOpts...)
|
||||
|
||||
// register services
|
||||
pb.RegisterUserServiceServer(srv, NewUserService())
|
||||
pb.RegisterOrderServiceServer(srv, NewOrderService())
|
||||
|
||||
// health check
|
||||
hs := health.NewServer()
|
||||
hs.SetServingStatus("", grpc_health_v1.HealthCheckResponse_SERVING)
|
||||
grpc_health_v1.RegisterHealthServer(srv, hs)
|
||||
|
||||
// reflection: dev only
|
||||
if os.Getenv("ENV") != "production" {
|
||||
reflection.Register(srv)
|
||||
}
|
||||
|
||||
return srv, nil
|
||||
}
|
||||
|
||||
func main() {
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
srv, err := NewGRPCServer()
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
lis, err := net.Listen("tcp", ":50051")
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
errCh := make(chan error, 1)
|
||||
go func() {
|
||||
errCh <- srv.Serve(lis)
|
||||
}()
|
||||
|
||||
select {
|
||||
case err := <-errCh:
|
||||
return
|
||||
case <-ctx.Done():
|
||||
// GracefulShutdown with timeout
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
srv.GracefulStop()
|
||||
close(done)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
log.Println("server stopped gracefully")
|
||||
case <-shutdownCtx.Done():
|
||||
log.Println("shutdown timeout, forcing stop")
|
||||
srv.Stop()
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
这个模板涵盖了:TLS、大小限制、keepalive、拦截器链、健康检查、reflection 条件开关、signal handling、graceful shutdown + 超时兜底。复制粘贴后即可投入生产使用。
|
||||
|
||||
## 关键概念对照表
|
||||
|
||||
| 概念 | 对应 API | 一句话总结 |
|
||||
|------|---------|-----------|
|
||||
| Server 创建 | `grpc.NewServer(...)` | 传入 ServerOption 配置行为 |
|
||||
| Service 注册 | `RegisterXxxServer(s, impl)` | 将实现绑定到 server dispatch table |
|
||||
| 零-stub 兼容 | 嵌入 `UnimplementedXxxServer` | proto 新增 method 编译期自动告警 |
|
||||
| 健康检查 | `health.NewServer()` | Kubernetes probe / SLA 监控对接 |
|
||||
| 调试反射 | `reflection.Register(s)` | 仅限开发环境 |
|
||||
| 优雅关闭 | `GracefulStop()` + 超时 | 等活跃请求完成,超时无情斩断 |
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[hhs/gRPC/3. 服务端实现/09-Streaming Handler]]
|
||||
- [[hhs/gRPC/3. 服务端实现/10-健康检查与反射]]
|
||||
- [[hhs/gRPC/1. 基础概念/02-gRPC 核心术语]]
|
||||
Reference in New Issue
Block a user