Skip to content

gRPC 客户端与连接管理

gRPC 客户端看似只需 grpc.Dial + 调用方法,但生产级的连接管理涉及超时、重试、负载均衡、连接状态监控等多个维度。本篇深入剖析 ClientConn 的内部机制,讲解 grpc.Dialgrpc.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.WithBlockgrpc.WithTimeoutgrpc.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 通常只配 UNAVAILABLEDEADLINE_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_robinleast_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.Builderresolver.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(推荐,非阻塞)vs grpc.Dial(旧,可阻塞),理解懒连接语义。
  • 连接选项:传输凭证、keepalive、消息大小、User-Agent。
  • 拦截器:一元和流两种,可链式组合,注入 request-id、日志、监控。
  • 超时控制:完全由 context 驱动,deadline 通过 HTTP/2 头部自动传播。
  • 重试机制:通过 ServiceConfig 声明,支持指数退避,仅对幂等操作和指定状态码生效。
  • 负载均衡:pick_first(默认)和 round_robin,通过 ServiceConfig 切换。
  • 自定义 Resolver:实现 resolver.Builder 接口,对接 etcd/Consul 等服务发现。
  • 状态监控GetState + WaitForStateChange 感知连接健康度。
  • 完整示例:服务发现 + 负载均衡 + 重试 + 降级的生产级客户端组合。

下一篇将聚焦拦截器与中间件模式,实现日志、认证、限流、指标采集等生产级横切逻辑。