Appearance
gRPC 客户端与连接管理
gRPC 客户端看似只需 grpc.Dial + 调用方法,但生产级的连接管理涉及超时、重试、负载均衡、连接状态监控等多个维度。本篇深入剖析 ClientConn 的内部机制,讲解 grpc.Dial 与 grpc.NewClient 的区别、keepalive 配置、客户端拦截器、服务配置重试策略、负载均衡策略与自定义 Resolver,帮助你构建健壮的 gRPC 客户端。
一、建立连接:grpc.Dial 和 grpc.NewClient
1. grpc.Dial 的历史与问题
grpc.Dial 是经典的客户端连接函数。它默认是非阻塞的——即使目标地址连不上,也会返回一个 ClientConn,真正的连接在后台异步建立。如果传了 grpc.WithBlock(),则会阻塞直到连接建立成功或超时。
go
package main
import (
"context"
"fmt"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// 经典 Dial 方式(已标记 deprecated,但仍广泛使用)
// 非阻塞:立即返回,连接在后台建立
conn, err := grpc.Dial("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("dial returned (non-blocking): %T\n", conn)
// 阻塞方式:WithBlock + WithTimeout
// 注意:grpc.WithTimeout 也已 deprecated,推荐用 context
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
conn2, err := grpc.DialContext(ctx, "127.0.0.1:50052",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithBlock(), // 阻塞直到连接建立
)
if err != nil {
fmt.Printf("dial 50052 (blocked): %v\n", err)
} else {
defer conn2.Close()
fmt.Println("dial 50052 connected")
}
}grpc.Dial 的问题在于:grpc.WithBlock、grpc.WithTimeout、grpc.WithReturnConnectionError 等选项与 context 的交互不够清晰,且 Dial 的语义("拨号")暗示同步建立连接,但默认行为是异步的,容易误导。
2. grpc.NewClient:推荐的新方式
从 gRPC-Go v1.63 开始,官方推荐使用 grpc.NewClient。它始终非阻塞,连接在后台懒加载建立,且不再接受 WithBlock/WithTimeout 等阻塞选项,超时完全由调用时的 context 控制。
go
package main
import (
"fmt"
"log"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// 推荐方式:grpc.NewClient(始终非阻塞,语义清晰)
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
// 这里的 error 仅表示选项配置有误,而非连接失败
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("NewClient returned: %T\n", conn)
// 连接是否真正建立,由首次 RPC 调用时的 context 超时决定
// 如果服务端不可达,第一次 RPC 会在 context 超时后返回 error
}何时用哪个:
- 新项目优先用
grpc.NewClient,超时由调用 context 控制。 - 维护老项目可继续用
grpc.DialContext+WithBlock。
二、连接选项详解
1. 传输凭证
go
package main
import (
"fmt"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// 无凭证(仅用于开发/测试,gRPC 1.63+ 推荐 insecure.NewCredentials())
// 旧代码中的 grpc.WithInsecure() 已废弃
conn1, _ := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
defer conn1.Close()
fmt.Println("insecure conn created")
// TLS 凭证
// creds, err := credentials.NewClientTLSFromFile("server.crt", "example.com")
// if err != nil { log.Fatal(err) }
// conn2, _ := grpc.NewClient("example.com:443",
// grpc.WithTransportCredentials(creds),
// )
// 使用自定义 tls.Config(如跳过证书校验,仅用于测试)
conn3, _ := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(credentials.NewTLS(nil)),
)
defer conn3.Close()
fmt.Println("TLS conn created")
}2. 连接池与 keepalive
gRPC 的 ClientConn 内部维护了一个 HTTP/2 连接池。默认情况下,一个 ClientConn 到一个目标地址只有一个底层连接,但这条连接上的多个 HTTP/2 stream 可以并发承载多个 RPC。
go
package main
import (
"fmt"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/keepalive"
)
func main() {
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
// 客户端 keepalive:定期检测连接存活
grpc.WithKeepaliveParams(keepalive.ClientParameters{
// 每 30 秒发送一次 ping
Time: 30 * time.Second,
// ping 超时 10 秒无响应则认为连接断开
Timeout: 10 * time.Second,
// 没有活动流时也允许 ping
PermitWithoutStream: true,
}),
// 单条 RPC 消息大小限制
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(16*1024*1024),
grpc.MaxCallSendMsgSize(16*1024*1024),
),
// 自定义 User-Agent
grpc.WithUserAgent("my-shop-client/1.0"),
)
if err != nil {
fmt.Printf("error: %v\n", err)
return
}
defer conn.Close()
fmt.Printf("conn with keepalive: %T\n", conn)
}连接复用要点:
ClientConn是线程安全的,可以被多个 goroutine 共享。- 不要每次 RPC 都 Dial——这是最常见的性能反模式。一个进程内一个目标地址只需一个
ClientConn。 - 如果需要更多并行度,可以创建多个
ClientConn,但通常单个连接的 HTTP/2 多路复用已足够。
三、客户端拦截器
1. 一元拦截器
go
package main
import (
"context"
"fmt"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
// 一元客户端拦截器签名
// func(ctx, method, req, reply, cc, invoker, opts...) error
func loggingClientInterceptor(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 requestIDClientInterceptor(ctx context.Context, method string,
req, reply interface{}, cc *grpc.ClientConn,
invoker grpc.UnaryInvoker, opts ...grpc.CallOption,
) error {
// 注入 request-id 到 metadata(简化演示)
fmt.Printf("[client] injecting request-id for %s\n", method)
return invoker(ctx, method, req, reply, cc, opts...)
}
func main() {
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithChainUnaryInterceptor(
requestIDClientInterceptor,
loggingClientInterceptor,
),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("client with interceptors: %T\n", conn)
}2. 流拦截器
go
package main
import (
"context"
"fmt"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
// 包装 ClientStream 以拦截收发
type wrappedClientStream struct {
grpc.ClientStream
method string
}
func (w *wrappedClientStream) RecvMsg(m interface{}) error {
err := w.ClientStream.RecvMsg(m)
if err == nil {
log.Printf("[client stream %s] recv msg", w.method)
}
return err
}
func (w *wrappedClientStream) SendMsg(m interface{}) error {
log.Printf("[client stream %s] send msg", w.method)
return w.ClientStream.SendMsg(m)
}
func loggingStreamClientInterceptor(ctx context.Context, desc *grpc.StreamDesc,
cc *grpc.ClientConn, method string, streamer grpc.Streamer,
opts ...grpc.CallOption,
) (grpc.ClientStream, error) {
start := time.Now()
s, err := streamer(ctx, desc, cc, method, opts...)
if err != nil {
return nil, err
}
return &wrappedClientStream{ClientStream: s, method: method}, func() {}() && nil || nil // simplified
}
// 正确实现
func loggingStreamClient(ctx context.Context, desc *grpc.StreamDesc,
cc *grpc.ClientConn, method string, streamer grpc.Streamer,
opts ...grpc.CallOption,
) (grpc.ClientStream, error) {
start := time.Now()
s, err := streamer(ctx, desc, cc, method, opts...)
if err != nil {
return nil, err
}
wrapped := &wrappedClientStream{ClientStream: s, method: method}
// 用 finalizer 记录总耗时(简化:记录到日志)
go func() {
<-ctx.Done()
log.Printf("[client stream %s] total cost=%v", method, time.Since(start))
}()
return wrapped, nil
}
func main() {
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithStreamInterceptor(loggingStreamClient),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("client with stream interceptor: %T\n", conn)
}四、超时控制
gRPC 客户端的超时完全依赖 context。每个 RPC 调用都应该带一个带超时的 context。
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"
)
// 模拟一个 RPC 调用
func fakeRPC(ctx context.Context) (string, error) {
select {
case <-time.After(2 * time.Second): // 模拟服务端处理 2 秒
return "done", nil
case <-ctx.Done():
return "", ctx.Err()
}
}
func main() {
conn, _ := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
defer conn.Close()
// 超时分级策略
callWithTimeout := func(name string, timeout time.Duration) {
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
start := time.Now()
_, err := fakeRPC(ctx)
cost := time.Since(start)
if err != nil {
st, _ := status.FromError(err)
fmt.Printf("%s: FAILED in %v, code=%v\n", name, cost, st.Code())
} else {
fmt.Printf("%s: OK in %v\n", name, cost)
}
}
callWithTimeout("fast-call (1s timeout)", 1*time.Second) // 超时
callWithTimeout("slow-call (3s timeout)", 3*time.Second) // 成功
// 传播超时:父 context 取消时子 context 也会取消
parentCtx, parentCancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer parentCancel()
childCtx, childCancel := context.WithTimeout(parentCtx, 5*time.Second) // 子超时更长
defer childCancel()
start := time.Now()
_, err := fakeRPC(childCtx)
fmt.Printf("child call: cost=%v err=%v (parent timeout wins)\n", time.Since(start), err)
_ = conn
_ = codes.DeadlineExceeded
_ = log.Printf
}超时层级传播是微服务的核心实践:网关设置总超时(如 5s),每经过一层服务减去已消耗时间作为下游的超时。gRPC 会自动通过 grpc-timeout HTTP/2 头部传播 deadline。
五、重试机制
gRPC 内置了重试机制,通过**服务配置(ServiceConfig)**声明。重试对透明故障转移非常重要,但需要服务端配合(必须是幂等操作)。
go
package main
import (
"context"
"fmt"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// 通过 ServiceConfig 配置重试策略
// methodConfig 匹配特定方法,retryPolicy 定义重试行为
retrySvcCfg := `{
"methodConfig": [{
"name": [
{"service": "shop.v1.OrderService", "method": "GetOrder"},
{"service": "shop.v1.OrderService", "method": "CreateOrder"}
],
"retryPolicy": {
"maxAttempts": 4,
"initialBackoff": "0.1s",
"maxBackoff": "1s",
"backoffMultiplier": 2.0,
"retryableStatusCodes": ["UNAVAILABLE", "DEADLINE_EXCEEDED"]
},
"timeout": "3s"
},{
"name": [
{"service": "shop.v1.OrderService"}
],
"timeout": "10s"
}]
}`
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(retrySvcCfg),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("client with retry config: %T\n", conn)
// 重试策略解析:
// maxAttempts: 4(1 次初始 + 3 次重试)
// initialBackoff: 第一次重试前等待 0.1s
// maxBackoff: 单次等待上限 1s
// backoffMultiplier: 指数退避乘数 2.0
// -> 0.1s, 0.2s, 0.4s, 0.8s, 1s, 1s...
// retryableStatusCodes: 只有这些状态码才重试
_ = context.Background
_ = time.Second
}重试注意事项:
- 只有幂等操作(GET、PUT)才适合重试,非幂等操作(POST/CREATE)重试可能导致重复创建。
retryableStatusCodes通常只配UNAVAILABLE和DEADLINE_EXCEEDED,不要配INTERNAL(可能是业务错误)。- gRPC 重试默认是 per-attempt transparent:客户端看到的是最终结果,中间重试过程不可见。
- 重试会消耗重试预算(retry budget),防止重试风暴。
六、负载均衡策略
gRPC 客户端内置了负载均衡器,支持多种 picker 策略。
1. pick_first(默认)
go
package main
import (
"fmt"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// pick_first:依次尝试所有地址,选中第一个能连上的,之后一直用这个
// 适合单实例或不需要负载均衡的场景
conn, _ := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(`{"loadBalancingConfig":[{"pick_first":{}}]}`),
)
defer conn.Close()
fmt.Printf("pick_first conn: %T\n", conn)
}2. round_robin
go
package main
import (
"fmt"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
// round_robin:在所有已连接的地址间轮询
// 需要后端地址列表(通常通过 DNS 或服务发现提供)
conn, _ := grpc.NewClient("dns:///shop-service:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(`{"loadBalancingConfig":[{"round_robin":{}}]}`),
)
defer conn.Close()
fmt.Printf("round_robin conn: %T\n", conn)
// DNS 地址格式说明:
// dns:///host:port -> 通过 DNS 解析 host 获取地址列表
// 127.0.0.1:50051 -> 单地址(pick_first 默认)
// passthrough:///127.0.0.1:50051 -> 不解析,直接使用
}两种策略对比:
| 策略 | 行为 | 适用场景 |
|---|---|---|
pick_first | 选第一个可达地址,故障时切换 | 单实例、主从切换 |
round_robin | 在所有可达地址间轮询 | 多实例负载均衡 |
gRPC 还支持 weighted_round_robin、least_request 等高级策略(通过 xDS 或自定义 balancer)。
七、自定义 Resolver
当后端地址不是通过 DNS 而是通过服务发现(如 etcd、Consul、Nacos)管理时,需要自定义 Resolver 来动态解析地址列表。
go
package main
import (
"context"
"fmt"
"log"
"sync"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/resolver"
)
// 自定义 Resolver:模拟从服务发现获取地址列表
type etcdResolverBuilder struct {
addrs []string // 模拟服务发现返回的地址
}
func (b *etcdResolverBuilder) Build(target resolver.Target, cc resolver.ClientConn,
opts resolver.BuildOptions,
) (resolver.Resolver, error) {
r := &etcdResolver{
cc: cc,
target: target,
current: b.addrs,
}
r.start() // 启动后台轮询
return r, nil
}
func (b *etcdResolverBuilder) Scheme() string {
return "etcd"
}
type etcdResolver struct {
cc resolver.ClientConn
target resolver.Target
mu sync.Mutex
current []string
}
func (r *etcdResolver) start() {
// 立即解析一次
r.resolve()
// 后台定期刷新(模拟服务发现的 watch 机制)
go func() {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for range ticker.C {
r.resolve()
}
}()
}
func (r *etcdResolver) resolve() {
r.mu.Lock()
defer r.mu.Unlock()
// 将地址列表转为 resolver.Address
addrs := make([]resolver.Address, len(r.current))
for i, a := range r.current {
addrs[i] = resolver.Address{Addr: a}
}
// 更新连接的状态和地址列表
r.cc.UpdateState(resolver.State{
Addresses: addrs,
// 可附加 ServiceConfig
})
}
func (r *etcdResolver) ResolveNow(o resolver.ResolveNowOptions) {
r.resolve()
}
func (r *etcdResolver) Close() {}
func main() {
// 注册自定义 Resolver
resolver.Register(&etcdResolverBuilder{
addrs: []string{
"127.0.0.1:50051",
"127.0.0.1:50052",
"127.0.0.1:50053",
},
})
// 使用自定义 scheme 的地址
conn, err := grpc.NewClient("etcd:///shop-service",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(`{"loadBalancingConfig":[{"round_robin":{}}]}`),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
fmt.Printf("conn with custom resolver: %T\n", conn)
_ = context.Background
}自定义 Resolver 的核心:实现 resolver.Builder 和 resolver.Resolver 接口,在 UpdateState 时传入最新的地址列表。配合 round_robin 策略,就实现了服务发现 + 客户端负载均衡。
八、连接状态监控
ClientConn 暴露了连接状态查询接口,可用于健康检查和故障感知。
go
package main
import (
"context"
"fmt"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/connectivity"
)
func main() {
conn, err := grpc.NewClient("127.0.0.1:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
// 获取当前连接状态
state := conn.GetState()
fmt.Printf("initial state: %s\n", state)
// 连接状态枚举:
// Idle - 空闲,未尝试连接
// Connecting - 正在连接
// Ready - 连接就绪,可以发 RPC
// TransientFailure - 临时失败(会自动重试)
// Shutdown - 连接已关闭
// 监控状态变化
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
go monitorState(ctx, conn)
// 触发连接(NewClient 是懒连接,第一次 RPC 才真正建立)
// 实际项目中这里会调用 RPC
time.Sleep(2 * time.Second)
fmt.Println("monitoring done")
}
func monitorState(ctx context.Context, conn *grpc.ClientConn) {
for {
state := conn.GetState()
fmt.Printf("[monitor] state=%s\n", state)
// WaitForStateChange 阻塞直到状态变化或 ctx 超时
if !conn.WaitForStateChange(ctx, state) {
// ctx 超时或取消
fmt.Println("[monitor] stopped (context done)")
return
}
select {
case <-ctx.Done():
return
default:
}
}
}
// 连接状态的典型转换:
// Idle -> Connecting -> Ready(成功)
// Idle -> Connecting -> TransientFailure -> Connecting -> Ready(重试成功)
// Ready -> Idle(连接断开,无活动 RPC)
// Ready -> Connecting -> TransientFailure(连接断开,有活动 RPC,持续重试)
func _unused() {
_ = connectivity.Ready
_ = connectivity.TransientFailure
}状态监控的使用场景:
- 健康检查面板:显示各下游服务的连接状态。
- 熔断触发:连续 TransientFailure 时触发熔断。
- 平滑启动:等待 Ready 状态再开始接收流量。
九、完整示例:带重试和负载均衡的客户端
go
package main
import (
"context"
"fmt"
"log"
"math/rand"
"os"
"os/signal"
"syscall"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/resolver"
"google.golang.org/grpc/status"
)
// === 自定义 Resolver:模拟服务发现 ===
type discoveryResolver struct {
cc resolver.ClientConn
addrs []string
stopCh chan struct{}
}
type discoveryBuilder struct {
initialAddrs []string
}
func (b *discoveryBuilder) Scheme() string { return "discovery" }
func (b *discoveryBuilder) Build(target resolver.Target, cc resolver.ClientConn,
opts resolver.BuildOptions,
) (resolver.Resolver, error) {
r := &discoveryResolver{
cc: cc,
addrs: b.initialAddrs,
stopCh: make(chan struct{}),
}
r.updateAddrs()
go r.watch()
return r, nil
}
func (r *discoveryResolver) updateAddrs() {
addrs := make([]resolver.Address, len(r.addrs))
for i, a := range r.addrs {
addrs[i] = resolver.Address{Addr: a}
}
r.cc.UpdateState(resolver.State{Addresses: addrs})
}
func (r *discoveryResolver) watch() {
ticker := time.NewTicker(15 * time.Second)
defer ticker.Stop()
for {
select {
case <-ticker.C:
// 模拟服务发现:偶尔增减实例
if len(r.addrs) > 1 && rand.Float32() < 0.3 {
r.addrs = r.addrs[:len(r.addrs)-1]
log.Printf("[discovery] instance removed, now %d instances", len(r.addrs))
r.updateAddrs()
}
case <-r.stopCh:
return
}
}
}
func (r *discoveryResolver) ResolveNow(resolver.ResolveNowOptions) {}
func (r *discoveryResolver) Close() { close(r.stopCh) }
// === 客户端拦截器 ===
func retryLoggingInterceptor(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...)
cost := time.Since(start)
if err != nil {
st, _ := status.FromError(err)
log.Printf("[rpc] %s FAILED code=%v cost=%v", method, st.Code(), cost)
} else {
log.Printf("[rpc] %s OK cost=%v", method, cost)
}
return err
}
// === 超时与降级封装 ===
func callWithFallback(ctx context.Context, method string,
rpcFunc func(context.Context) (string, error),
) (string, error) {
ctx, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel()
result, err := rpcFunc(ctx)
if err != nil {
st, _ := status.FromError(err)
switch st.Code() {
case codes.Unavailable:
return "fallback: service unavailable", nil
case codes.DeadlineExceeded:
return "fallback: timeout", nil
default:
return "", err
}
}
return result, nil
}
func main() {
// 注册服务发现 Resolver
resolver.Register(&discoveryBuilder{
initialAddrs: []string{
"127.0.0.1:50051",
"127.0.0.1:50052",
"127.0.0.1:50053",
},
})
// 构建带重试和负载均衡的连接
svcCfg := `{
"loadBalancingConfig": [{"round_robin": {}}],
"methodConfig": [{
"name": [{"service": "shop.v1.OrderService"}],
"retryPolicy": {
"maxAttempts": 3,
"initialBackoff": "0.2s",
"maxBackoff": "2s",
"backoffMultiplier": 2.0,
"retryableStatusCodes": ["UNAVAILABLE"]
},
"timeout": "3s"
}]
}`
conn, err := grpc.NewClient("discovery:///shop-service",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(svcCfg),
grpc.WithChainUnaryInterceptor(retryLoggingInterceptor),
grpc.WithKeepaliveParams(grpc.KeepaliveParams{}),
)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
// 等待退出信号
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
// 模拟周期性调用
go func() {
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
callCount := 0
for {
select {
case <-ticker.C:
callCount++
// 模拟 RPC 调用(真实场景用 protoc 生成的 client)
result, err := callWithFallback(context.Background(),
"GetOrder",
func(ctx context.Context) (string, error) {
// 这里模拟调用失败(因为服务端没起)
return "", status.Error(codes.Unavailable, "no server available")
},
)
fmt.Printf("call #%d: result=%q err=%v\n", callCount, result, err)
case <-quit:
return
}
}
}()
<-quit
fmt.Println("\nshutting down client...")
}这个示例综合了自定义 Resolver(服务发现)+ round_robin 负载均衡 + 内置重试 + 拦截器日志 + 超时降级,是生产级 gRPC 客户端的标配组合。
十、小结
本篇深入 gRPC 客户端的连接管理:
- 连接建立:
grpc.NewClient(推荐,非阻塞)vsgrpc.Dial(旧,可阻塞),理解懒连接语义。 - 连接选项:传输凭证、keepalive、消息大小、User-Agent。
- 拦截器:一元和流两种,可链式组合,注入 request-id、日志、监控。
- 超时控制:完全由 context 驱动,deadline 通过 HTTP/2 头部自动传播。
- 重试机制:通过 ServiceConfig 声明,支持指数退避,仅对幂等操作和指定状态码生效。
- 负载均衡:pick_first(默认)和 round_robin,通过 ServiceConfig 切换。
- 自定义 Resolver:实现
resolver.Builder接口,对接 etcd/Consul 等服务发现。 - 状态监控:
GetState+WaitForStateChange感知连接健康度。 - 完整示例:服务发现 + 负载均衡 + 重试 + 降级的生产级客户端组合。
下一篇将聚焦拦截器与中间件模式,实现日志、认证、限流、指标采集等生产级横切逻辑。