四、负载均衡、路由与集群

2021-02-18T14:21:02+08:00 | 32分钟阅读 | 更新于 2021-02-18T14:21:02+08:00

@

学习目标

学完本章你应该能够:

  1. 用自己的话讲清两类算法:区分"不实时计算负载"与"尝试实时计算负载"两大类负载均衡算法,说清它们各自依赖的核心假设与适用边界。
  2. 逐算法讲清取舍:对轮询、加权轮询、随机、加权随机、哈希、一致性哈希,能讲出各自的优缺点,并能用"假设不成立"解释生产上"某台机器被打爆"的现象。
  3. 在 gRPC 里自定义负载均衡器:理解并实现 PickerBuilder + Picker 两个接口,把负载均衡算法真正接进 gRPC 调用链。
  4. 设计路由策略:讲清标签路由、健康路由、分组(A/B)路由的思路,并理解"路由是负载均衡之前的一个过滤步骤"这个关系。
  5. 把 Cluster 讲成面试故事:说清 failover / failfast / 广播 / 组播的区别与使用场景,面试时能结合重试、换节点、网关等讲成一段完整叙事。

前置知识

  • 上一章《服务注册与发现》:知道"有哪些可用实例"是怎么来的(本章是在此基础上的下一步)。
  • Go 并发基础:sync.Mutexsync/atomic 的基本用法。
  • gRPC 基本用法:知道一次 RPC 调用是怎么发起的(不要求写过拦截器)。

本章你会动手做的事

  1. 跑一遍 RoundRobinPicker.Pick(),观察计数器取模如何循环选中实例。
  2. 手算一个 key 的哈希值,在一致性哈希环上"顺时针找最近节点",验证扩缩容只影响部分请求。
  3. 仿照分组路由示例,写一个按 A/B 标记把请求分流到不同实例组的 Picker

前言:这么多可用实例,我该把请求发给谁?

在前一章中,我们讨论了服务注册与发现,解决了"有哪些可用服务实例"的问题。现在我们面临下一个问题:

这么多可用服务实例,我该把请求发给谁?

这一步,我们直觉就想到关键字:负载均衡

从理论上来说,我们希望将请求发给那个能最快返回响应的实例。但现实中,我们很难实时知道每个实例的当前状态,于是就诞生了各种各样的负载均衡算法。

本教程涵盖以下主题:

  1. 负载均衡算法概览——两大类算法的分类与对比
  2. 不实时计算的负载均衡算法——轮询、加权轮询、随机、加权随机、哈希、一致性哈希
  3. 实时计算的负载均衡算法——最小连接数、最少活跃数、最快响应时间
  4. gRPC 负载均衡实现——PickerBuilder、Picker 接口与代码实现
  5. 路由策略——标签路由、健康路由、分组功能实现
  6. 集群抽象(Cluster)——failover、failfast、广播、组播
  7. 总结与面试要点

一、负载均衡算法概览

负载均衡的核心目标是:将请求发给那个能最快返回响应的实例。所有算法都在回答同一个问题——怎么找出这个实例?

目前主流的算法分为两大类:

类别算法特点
不实时计算负载轮询、加权轮询、随机、加权随机、哈希、一致性哈希不关心实例当前的实际负载,依赖统计规律
尝试实时计算负载最快响应时间、最少连接数、最少活跃数尝试获取实例当前的负载指标来做决策

理想情况下,我们希望能够实时问一下各个实例的负载。但实际上,实时获取负载本身就有开销,而且获取到的数据可能已经过时了。所以大多数时候,我们使用不实时计算的算法。

每种算法都有一些假设(assumption),理解这些假设是理解算法优缺点的关键。后续我们会逐一分析。


二、不实时计算的负载均衡算法

2.1 轮询(Round Robin)

原理:排排坐分果果——按顺序依次将请求分发给每个实例。

flowchart LR
    R[请求1] --> A[实例A]
    R2[请求2] --> B[实例B]
    R3[请求3] --> C[实例C]
    R4[请求4] --> A2[实例A]
    R5[请求5] --> B2[实例B]

假设

  • 所有服务器的处理能力是一样的
  • 所有请求所需的资源也是一样的

为什么大多数时候它运作效果都很好? 因为在请求量大的情况下,即使假设不完全成立,统计上请求会均匀分布到各个实例,整体效果趋于均衡。

轮询算法的 Go 实现

package loadbalance

import (
	"sync"
)

// ============================================================
// 轮询(Round Robin)负载均衡算法
// ============================================================

// RoundRobinPicker 是轮询负载均衡器
// 核心思想:维护一个计数器,每次请求时计数器 +1,
// 然后对实例总数取模,得到目标实例的索引
type RoundRobinPicker struct {
	mu        sync.Mutex    // 互斥锁,保证并发安全
	instances []string      // 可用实例列表,存储实例地址如 "172.10.0.1:9091"
	index     uint64        // 当前轮询到的位置
}

// NewRoundRobinPicker 创建一个轮询负载均衡器
// 参数 instances 是初始的可用实例列表
func NewRoundRobinPicker(instances []string) *RoundRobinPicker {
	return &RoundRobinPicker{
		instances: instances,
	}
}

// Pick 从可用实例中选择一个返回
// 每次调用返回下一个实例(按顺序循环)
func (r *RoundRobinPicker) Pick() string {
	r.mu.Lock()         // 加锁,防止并发时计数器混乱
	defer r.mu.Unlock() // 函数结束时释放锁

	if len(r.instances) == 0 {
		return "" // 没有可用实例,返回空字符串
	}

	// 取模运算得到当前索引
	// 例如有 3 个实例,index 依次为 0, 1, 2, 0, 1, 2, ...
	idx := r.index % uint64(len(r.instances))
	r.index++ // 计数器递增

	return r.instances[idx]
}

// Update 更新可用实例列表
// 当服务发现有变更时(新增或下线实例),调用此方法更新
func (r *RoundRobinPicker) Update(instances []string) {
	r.mu.Lock()
	defer r.mu.Unlock()
	r.instances = instances
}

2.2 加权轮询(Weighted Round Robin)

原理:也是排排坐分果果,但是按照权重来分——权重大的多分,权重少的少分。

假设

  • 用权重来代表服务器的处理能力(如 8 核机器权重为 8,4 核机器权重为 4)
  • 所有请求所需的资源也是一样的

Nginx 平滑加权轮询算法

设想一下,如果有三个实例,权重分别是 5、1、1,那么简单的加权轮询会让权重为 5 的实例连续被选中 5 次,造成瞬间压力集中。

平滑加权轮询通过动态调整权重来实现平滑效果。算法有三个值:

字段说明
weight固定权重,初始配置的权重值
currentWeight当前权重,每次选择时会变化
effectiveWeight有效权重,可根据调用结果动态调整(如出错时降低)

算法步骤

  1. 每次挑选实例时,计算所有实例的 effectiveWeight 之和作为 totalWeight
  2. 对于每一个实例,更新 currentWeight = currentWeight + effectiveWeight
  3. 挑选 currentWeight 最大的那个节点作为最终节点
  4. 更新被选中节点的 currentWeight = currentWeight - totalWeight

这个算法的精妙之处在于:权重大的节点不会被连续选中,而是分散在各个位置,从而达到"平滑"的效果。

平滑加权轮询的 Go 实现

package loadbalance

import (
	"sync"
)

// ============================================================
// 平滑加权轮询(Smooth Weighted Round Robin)负载均衡算法
// 参考 Nginx 的平滑加权轮询实现
// ============================================================

// WeightedInstance 代表一个带权重的实例
type WeightedInstance struct {
	Address           string // 实例地址,如 "172.10.0.1:9091"
	Weight            int    // 固定权重,用户配置的初始权重
	CurrentWeight     int    // 当前权重,算法运行过程中动态变化
	EffectiveWeight   int    // 有效权重,可根据调用结果动态调整
}

// SmoothWeightedRRPicker 是平滑加权轮询负载均衡器
type SmoothWeightedRRPicker struct {
	mu        sync.Mutex
	instances []*WeightedInstance // 所有可用实例
}

// NewSmoothWeightedRRPicker 创建一个平滑加权轮询负载均衡器
// 参数 instances 是带权重的实例列表
func NewSmoothWeightedRRPicker(instances []*WeightedInstance) *SmoothWeightedRRPicker {
	// 初始化:effectiveWeight 初始等于 weight
	for _, inst := range instances {
		inst.EffectiveWeight = inst.Weight
		inst.CurrentWeight = 0
	}
	return &SmoothWeightedRRPicker{
		instances: instances,
	}
}

// Pick 从可用实例中选择一个返回
// 核心算法:每次选择 currentWeight 最大的实例
func (p *SmoothWeightedRRPicker) Pick() string {
	p.mu.Lock()
	defer p.mu.Unlock()

	if len(p.instances) == 0 {
		return ""
	}

	// 第一步:计算 totalWeight(所有实例的 effectiveWeight 之和)
	// 同时更新每个实例的 currentWeight
	totalWeight := 0
	bestInstance := p.instances[0] // 记录 currentWeight 最大的实例

	for _, inst := range p.instances {
		// 更新当前权重:currentWeight += effectiveWeight
		inst.CurrentWeight += inst.EffectiveWeight
		// 累加总权重
		totalWeight += inst.EffectiveWeight

		// 找到 currentWeight 最大的实例
		if inst.CurrentWeight > bestInstance.CurrentWeight {
			bestInstance = inst
		}
	}

	// 第二步:被选中的实例,currentWeight 减去 totalWeight
	// 这一步是关键——它确保权重大的实例不会被连续选中
	// 减去 totalWeight 后,下次选择时该实例的 currentWeight 会大幅降低
	bestInstance.CurrentWeight -= totalWeight

	return bestInstance.Address
}

// OnSuccess 调用成功时调用,恢复有效权重
func (p *SmoothWeightedRRPicker) OnSuccess(addr string) {
	p.mu.Lock()
	defer p.mu.Unlock()
	for _, inst := range p.instances {
		if inst.Address == addr {
			// 调用成功,恢复有效权重(但不超过初始权重)
			if inst.EffectiveWeight < inst.Weight {
				inst.EffectiveWeight++
			}
			break
		}
	}
}

// OnError 调用失败时调用,降低有效权重
func (p *SmoothWeightedRRPicker) OnError(addr string) {
	p.mu.Lock()
	defer p.mu.Unlock()
	for _, inst := range p.instances {
		if inst.Address == addr {
			// 调用失败,降低有效权重
			// 这样下次选择时,该实例被选中的概率会降低
			inst.EffectiveWeight--
			// 注意:要防止 effectiveWeight 变成 0 或负数
			// 否则该实例将永远不被选中
			if inst.EffectiveWeight < 1 {
				inst.EffectiveWeight = 1
			}
			break
		}
	}
}

权重调整的注意事项:调整权重并不是说一定是 +1 -1,要考虑自己权重的实际取值。一个实例的权重如果是动态调整的,那么就要考虑上限和下限的问题,尤其要考虑调整权重的过程中会不会导致权重变成 0、最大值或者最小值——这三个值可能导致实例完全不被选中,或者一直被选中。

2.3 随机(Random)

原理:闭着眼睛瞎选——每次请求随机选一个实例。

假设(在轮询的两个假设之上,还多一个):

  • 所有服务器的处理能力是一样的
  • 所有请求所需的资源也是一样的
  • 每台服务器被随机到的概率是一样的

为什么随机也能工作? 在大量请求的情况下,根据概率论的大数定律,服务器之间的负载会趋于一样。相比之下,轮询的可控性更强,但大多数时候可以认为它们效果差不多。

package loadbalance

import (
	"math/rand"
	"sync"
	"time"
)

// ============================================================
// 随机(Random)负载均衡算法
// ============================================================

// RandomPicker 是随机负载均衡器
type RandomPicker struct {
	mu        sync.Mutex
	instances []string
	rand      *rand.Rand // 随机数生成器
}

// NewRandomPicker 创建一个随机负载均衡器
func NewRandomPicker(instances []string) *RandomPicker {
	return &RandomPicker{
		instances: instances,
		// 使用当前时间作为随机种子,确保每次启动结果不同
		rand: rand.New(rand.NewSource(time.Now().UnixNano())),
	}
}

// Pick 随机选择一个实例
func (r *RandomPicker) Pick() string {
	r.mu.Lock()
	defer r.mu.Unlock()

	if len(r.instances) == 0 {
		return ""
	}

	// 生成一个 [0, len) 范围内的随机索引
	idx := r.rand.Intn(len(r.instances))
	return r.instances[idx]
}

2.4 加权随机(Weighted Random)

原理:根据权重来确定选中概率。例如三个实例权重为 5、3、2,总共 10,那么第一个实例有 50% 的概率被选中。

假设

  • 用权重来代表服务器的处理能力
  • 所有请求所需的资源也是一样的
package loadbalance

import (
	"math/rand"
	"sync"
	"time"
)

// ============================================================
// 加权随机(Weighted Random)负载均衡算法
// ============================================================

// WeightedRandomPicker 是加权随机负载均衡器
type WeightedRandomPicker struct {
	mu        sync.Mutex
	instances []*WeightedInstance
	rand      *rand.Rand
}

// NewWeightedRandomPicker 创建一个加权随机负载均衡器
func NewWeightedRandomPicker(instances []*WeightedInstance) *WeightedRandomPicker {
	return &WeightedRandomPicker{
		instances: instances,
		rand:      rand.New(rand.NewSource(time.Now().UnixNano())),
	}
}

// Pick 根据权重随机选择一个实例
// 原理:将权重看作区间长度,随机数落在哪个区间就选哪个实例
// 例如权重为 [5, 3, 2],总权重为 10
// 实例1 占 [0, 5),实例2 占 [5, 8),实例3 占 [8, 10)
func (p *WeightedRandomPicker) Pick() string {
	p.mu.Lock()
	defer p.mu.Unlock()

	if len(p.instances) == 0 {
		return ""
	}

	// 第一步:计算总权重
	totalWeight := 0
	for _, inst := range p.instances {
		totalWeight += inst.Weight
	}

	// 第二步:生成一个 [0, totalWeight) 范围内的随机数
	r := p.rand.Intn(totalWeight)

	// 第三步:遍历实例,找到随机数落在的区间
	acc := 0 // 累加权重
	for _, inst := range p.instances {
		acc += inst.Weight
		if r < acc {
			return inst.Address
		}
	}

	// 理论上不会走到这里
	return p.instances[len(p.instances)-1].Address
}

2.5 哈希(Hash)

原理:对请求的某个特征(如用户 ID)计算哈希值,然后对实例数量取模,决定发送到哪个实例。

假设

  • 所有服务器的处理能力是一样的
  • 所有请求所需的资源也是一样的
  • 哈希值是均匀的

缺点:哈希值不均匀会导致请求堆积在一个地方。另外,当实例数量变化时(如扩容/缩容),大部分请求的哈希结果都会改变,导致缓存失效等问题。

package loadbalance

import (
	"hash/fnv"
	"sync"
)

// ============================================================
// 哈希(Hash)负载均衡算法
// ============================================================

// HashPicker 是哈希负载均衡器
// 核心思想:对请求的 key 计算 hash,然后对实例数取模
type HashPicker struct {
	mu        sync.Mutex
	instances []string
}

// NewHashPicker 创建一个哈希负载均衡器
func NewHashPicker(instances []string) *HashPicker {
	return &HashPicker{
		instances: instances,
	}
}

// Pick 根据 key 选择一个实例
// 参数 key 可以是用户 ID、请求 ID 等业务标识
// 相同的 key 总是路由到同一个实例(前提是实例列表不变)
func (p *HashPicker) Pick(key string) string {
	p.mu.Lock()
	defer p.mu.Unlock()

	if len(p.instances) == 0 {
		return ""
	}

	// 第一步:使用 FNV 哈希算法计算 key 的哈希值
	h := fnv.New32a()
	h.Write([]byte(key))
	hashValue := h.Sum32()

	// 第二步:对实例数量取模,得到目标实例索引
	idx := hashValue % uint32(len(p.instances))

	return p.instances[idx]
}

2.6 一致性哈希(Consistent Hash)

原理:一致性哈希是哈希算法的改进,引入了一个环状的思路。同样是计算哈希值,但是哈希值和节点的关系是:哈希值落在某一个区间内的,会被特定的一个节点处理。

flowchart TB
    subgraph 一致性哈希环
        direction TB
        N0["0"] --- N1["节点A
(hash=1000)"] N1 --- N2["节点B
(hash=3000)"] N2 --- N3["节点C
(hash=7000)"] N3 --- N4["节点D
(hash=9000)"] N4 --- N5["2^32-1"] end K1["请求key
hash=1500"] -.->|顺时针找最近节点| N2 K2["请求key
hash=5000"] -.->|顺时针找最近节点| N3 K3["请求key
hash=8500"] -.->|顺时针找最近节点| N4

优点:当增加节点或减少节点的时候,只有一部分请求命中的节点会发生变化,而不会像普通哈希那样大部分请求都要重新分配。

虚拟节点:为了解决节点数量少时哈希分布不均匀的问题,一致性哈希引入了虚拟节点——每个物理节点对应多个虚拟节点,虚拟节点均匀分布在环上。

package loadbalance

import (
	"hash/fnv"
	"sort"
	"sync"
)

// ============================================================
// 一致性哈希(Consistent Hash)负载均衡算法
// ============================================================

// ConsistentHashPicker 是一致性哈希负载均衡器
// 核心思想:将节点和请求都映射到一个 0 ~ 2^32-1 的环上
// 请求顺时针找到的第一个节点就是目标节点
type ConsistentHashPicker struct {
	mu           sync.Mutex
	replicas     int            // 每个物理节点的虚拟节点数量
	ring         []uint32       // 哈希环,存储所有虚拟节点的哈希值(有序)
	hashMap      map[uint32]string // 哈希值 -> 物理节点地址的映射
}

// NewConsistentHashPicker 创建一个一致性哈希负载均衡器
// 参数 replicas 是每个物理节点的虚拟节点数,通常设为 100~200
// 虚拟节点越多,分布越均匀,但内存开销也越大
func NewConsistentHashPicker(replicas int) *ConsistentHashPicker {
	return &ConsistentHashPicker{
		replicas: replicas,
		hashMap:  make(map[uint32]string),
	}
}

// Add 添加物理节点到哈希环
// 每个物理节点会创建 replicas 个虚拟节点
func (p *ConsistentHashPicker) Add(nodes ...string) {
	p.mu.Lock()
	defer p.mu.Unlock()

	for _, node := range nodes {
		// 为每个物理节点创建 replicas 个虚拟节点
		for i := 0; i < p.replicas; i++ {
			// 虚拟节点的 key = 物理节点地址 + 编号
			// 例如 "172.10.0.1:9091#0", "172.10.0.1:9091#1", ...
			virtualKey := node + "#" + string(rune(i))
			hash := p.hashKey(virtualKey)

			// 将虚拟节点的哈希值加入环
			p.ring = append(p.ring, hash)
			// 记录哈希值对应的物理节点
			p.hashMap[hash] = node
		}
	}

	// 保持环有序,方便二分查找
	sort.Slice(p.ring, func(i, j int) bool {
		return p.ring[i] < p.ring[j]
	})
}

// Pick 根据 key 选择一个实例
// 原理:计算 key 的哈希值,在环上顺时针找到第一个节点
func (p *ConsistentHashPicker) Pick(key string) string {
	p.mu.Lock()
	defer p.mu.Unlock()

	if len(p.ring) == 0 {
		return ""
	}

	// 第一步:计算 key 的哈希值
	hash := p.hashKey(key)

	// 第二步:在环上二分查找第一个 >= hash 的虚拟节点
// sort.Search 返回第一个满足条件的索引
	idx := sort.Search(len(p.ring), func(i int) bool {
		return p.ring[i] >= hash
	})

	// 如果所有虚拟节点的哈希值都小于 hash,则回到环的起点(idx = 0)
	if idx == len(p.ring) {
		idx = 0
	}

	// 返回虚拟节点对应的物理节点地址
	return p.hashMap[p.ring[idx]]
}

// hashKey 计算字符串的哈希值
func (p *ConsistentHashPicker) hashKey(key string) uint32 {
	h := fnv.New32a()
	h.Write([]byte(key))
	return h.Sum32()
}

三、实时计算的负载均衡算法

这类算法尝试实时获取实例的负载指标来做决策,理论上比不实时计算的算法更精准,但获取指标本身也有开销。

3.1 最小连接数(Least Connections)

假设

  • 用连接数来代表服务器负载
  • 所有服务器的处理能力是一样的
  • 请求所需资源都一样

缺点

  • 可能会短时间内把所有的请求都发到同一台服务器上
  • 连接复用的情况下,连接数不能很好地代表服务器的负载——比如 gRPC 使用 HTTP/2 多路复用,一个连接上可以并发多个请求

注意:在 gRPC 里,因为我们不直接管理连接(gRPC 自己管理连接池),所以这个算法是实现不了的。

3.2 最少活跃数(Least Active)

假设

  • 用服务器上的请求数量来代表负载
  • 所有服务器的处理能力是一样的
  • 请求所需资源都一样

原理:活跃数 = 已发送但还未收到响应的请求数量。选择活跃数最少的实例。

缺点:同样可能短时间内把所有请求发到同一台服务器上。

3.3 最快响应时间(Fastest Response)

假设

  • 用服务器的响应时间来代表负载
  • 请求所需资源都一样

缺点:偶尔几个慢请求会导致节点接下来不太可能会被选上。极端情况下,节点永远不会被选上。

package loadbalance

import (
	"sync"
	"time"
)

// ============================================================
// 最少活跃数(Least Active)负载均衡算法
// ============================================================

// ActiveInstance 记录实例的活跃请求数
type ActiveInstance struct {
	Address      string
	ActiveCount  int64       // 当前活跃请求数(已发送但未收到响应)
	TotalTime    time.Duration // 总响应时间(用于计算平均响应时间)
	TotalCount   int64       // 总请求数
}

// LeastActivePicker 是最少活跃数负载均衡器
type LeastActivePicker struct {
	mu        sync.Mutex
	instances map[string]*ActiveInstance // 实例地址 -> 活跃信息
}

// NewLeastActivePicker 创建一个最少活跃数负载均衡器
func NewLeastActivePicker(addrs []string) *LeastActivePicker {
	m := make(map[string]*ActiveInstance)
	for _, addr := range addrs {
		m[addr] = &ActiveInstance{Address: addr}
	}
	return &LeastActivePicker{instances: m}
}

// Pick 选择活跃数最少的实例
func (p *LeastActivePicker) Pick() string {
	p.mu.Lock()
	defer p.mu.Unlock()

	var bestAddr string
	var minActive int64 = -1

	for addr, inst := range p.instances {
		// 找到活跃数最少的实例
		if minActive == -1 || inst.ActiveCount < minActive {
			minActive = inst.ActiveCount
			bestAddr = addr
		}
	}

	// 选中后,活跃数 +1
	if bestAddr != "" {
		p.instances[bestAddr].ActiveCount++
	}

	return bestAddr
}

// OnResponse 收到响应时调用
// 记录响应时间,活跃数 -1
func (p *LeastActivePicker) OnResponse(addr string, duration time.Duration) {
	p.mu.Lock()
	defer p.mu.Unlock()

	if inst, ok := p.instances[addr]; ok {
		inst.ActiveCount--  // 活跃数减一
		inst.TotalTime += duration
		inst.TotalCount++
	}
}

四、gRPC 负载均衡实现

4.1 gRPC 负载均衡架构

在 gRPC 里面要自定义负载均衡算法,需要实现两个接口:

接口作用
balancer.PickerBuilder构建 Picker,接收可用连接列表
balancer.Picker执行负载均衡,每次请求时选择一个连接

这非常类似于注册中心的设计——也是实现两个接口。不同的是,负载均衡是通过 ServiceConfig 来指定的。

flowchart TB
    subgraph gRPC 负载均衡流程
        R[Resolver
服务发现] --> B[Balancer
负载均衡器] B --> PB[PickerBuilder
构建Picker] PB --> P[Picker
选择连接] P --> C1[连接1] P --> C2[连接2] P --> C3[连接3] end

4.2 PickerBuilder 和 Picker 接口

package loadbalance

import (
	"google.golang.org/grpc/balancer"
	"google.golang.org/grpc/resolver"
)

// ============================================================
// gRPC 负载均衡接口定义
// ============================================================

// PickerBuilder 接口:构建 Picker
// gRPC 在连接信息变更时会调用 Build 方法,传入最新的可用连接列表
// 我们需要返回一个 Picker 实例
type PickerBuilder interface {
	// Build 构建 Picker
	// 参数 readySCs 是所有就绪的子连接(SubConn),每个 SubConn 对应一个服务实例
	Build(readySCs map[balancer.SubConn]balancer.SubConnInfo) balancer.Picker
}

// Picker 接口:执行负载均衡
// 每次 RPC 调用时,gRPC 会调用 Pick 方法来选择一个连接
type Picker interface {
	// Pick 选择一个子连接
	// 参数 info 包含本次调用的上下文信息(如 metadata)
	// 返回选中的 SubConn 和是否完成
	Pick(info balancer.PickInfo) (balancer.PickResult, error)
}

4.3 轮询负载均衡实现

下面是在 gRPC 中实现轮询负载均衡的完整代码:

package loadbalance

import (
	"sync"
	"sync/atomic"

	"google.golang.org/grpc/balancer"
	"google.golang.org/grpc/balancer/base"
	"google.golang.org/grpc/resolver"
)

// ============================================================
// gRPC 轮询负载均衡完整实现
// ============================================================

// roundRobinPicker gRPC 轮询 Picker 实现
type roundRobinPicker struct {
	subConns []balancer.SubConn // 所有可用的子连接
	mu       sync.Mutex
	index    uint64 // 当前轮询位置(使用原子操作保证并发安全)
}

// Pick 实现 Picker 接口
// 每次 RPC 调用时,gRPC 会调用此方法选择一个连接
func (p *roundRobinPicker) Pick(info balancer.PickInfo) (balancer.PickResult, error) {
	if len(p.subConns) == 0 {
		// 没有可用连接,返回错误
		return balancer.PickResult{}, balancer.ErrNoSubConnAvailable
	}

	// 原子递增索引,确保并发安全
	idx := atomic.AddUint64(&p.index, 1)
	// 取模得到当前连接的索引
	subConn := p.subConns[idx%uint64(len(p.subConns))]

	// 返回选中的连接
	// PickResult 包含 SubConn 和可选的 Done 回调函数
	return balancer.PickResult{
		SubConn: subConn,
	}, nil
}

// roundRobinPickerBuilder gRPC 轮询 PickerBuilder 实现
type roundRobinPickerBuilder struct{}

// Build 实现 PickerBuilder 接口
// gRPC 在连接列表变更时调用此方法,传入最新的可用连接
func (b *roundRobinPickerBuilder) Build(readySCs map[balancer.SubConn]balancer.SubConnInfo) balancer.Picker {
	// 将 map 转换为切片,方便轮询
	var subConns []balancer.SubConn
	for sc := range readySCs {
		subConns = append(subConns, sc)
	}

	return &roundRobinPicker{
		subConns: subConns,
	}
}

// ============================================================
// 注册自定义负载均衡器到 gRPC
// ============================================================

// BalancerName 负载均衡器名称
const BalancerName = "my_round_robin"

// RegisterBalancer 注册自定义负载均衡器
// 注册后,可以在创建 gRPC 客户端时通过 ServiceConfig 指定使用此负载均衡器
func RegisterBalancer() {
	// base.NewBalancerBuilder 创建一个基础负载均衡器
	// 第一个参数是负载均衡器名称
	// 第二个参数是 PickerBuilder
	// 第三个参数是配置选项
	base.NewBalancerBuilder(BalancerName, &roundRobinPickerBuilder{}, base.Config{})
}

// ============================================================
// 客户端使用示例
// ============================================================

// 客户端创建时指定负载均衡器的方式:
//
// import (
//     "google.golang.org/grpc"
//     "google.golang.org/grpc/balancer/roundrobin"
//     "google.golang.org/grpc/resolver"
// )
//
// func main() {
//     // 方式一:使用 gRPC 内置的 roundrobin
//     conn, err := grpc.Dial(
//         "etcd:///user-service",  // 使用 etcd 服务发现
//         grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
//     )
//
//     // 方式二:使用自定义负载均衡器
//     RegisterBalancer() // 先注册
//     conn, err := grpc.Dial(
//         "etcd:///user-service",
//         grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"my_round_robin"}`),
//     )
// }

4.4 平滑加权轮询实现

package loadbalance

import (
	"sync"

	"google.golang.org/grpc/balancer"
	"google.golang.org/grpc/balancer/base"
	"google.golang.org/grpc/metadata"
	"google.golang.org/grpc/resolver"
)

// ============================================================
// gRPC 平滑加权轮询负载均衡实现
// ============================================================

// weightedSubConn 包装 SubConn 和其权重信息
type weightedSubConn struct {
	subConn         balancer.SubConn // gRPC 子连接
	weight          int              // 固定权重
	currentWeight   int              // 当前权重(动态变化)
	effectiveWeight int              // 有效权重(可动态调整)
	addr            string           // 实例地址,用于从 metadata 获取权重
}

// smoothWeightedPicker gRPC 平滑加权轮询 Picker
type smoothWeightedPicker struct {
	mu       sync.Mutex
	subConns []*weightedSubConn
}

// Pick 实现 Picker 接口
func (p *smoothWeightedPicker) Pick(info balancer.PickInfo) (balancer.PickResult, error) {
	if len(p.subConns) == 0 {
		return balancer.PickResult{}, balancer.ErrNoSubConnAvailable
	}

	p.mu.Lock()
	defer p.mu.Unlock()

	// 第一步:计算总权重,更新每个实例的 currentWeight
	totalWeight := 0
	best := p.subConns[0]

	for _, sc := range p.subConns {
		sc.currentWeight += sc.effectiveWeight
		totalWeight += sc.effectiveWeight

		if sc.currentWeight > best.currentWeight {
			best = sc
		}
	}

	// 第二步:被选中的实例 currentWeight 减去 totalWeight
	best.currentWeight -= totalWeight

	// 第三步:返回结果,设置 Done 回调用于动态调整权重
	return balancer.PickResult{
		SubConn: best.subConn,
		Done: func(info balancer.PickDoneInfo) {
			p.mu.Lock()
			defer p.mu.Unlock()

			if info.Err != nil {
				// 调用失败,降低有效权重
				best.effectiveWeight--
				if best.effectiveWeight < 1 {
					best.effectiveWeight = 1 // 防止降到 0
				}
			} else {
				// 调用成功,恢复有效权重
				if best.effectiveWeight < best.weight {
					best.effectiveWeight++
				}
			}
		},
	}, nil
}

// smoothWeightedBuilder gRPC 平滑加权轮询 Builder
type smoothWeightedBuilder struct{}

// Build 实现 PickerBuilder 接口
// 从 SubConnInfo 中提取权重信息
func (b *smoothWeightedBuilder) Build(readySCs map[balancer.SubConn]balancer.SubConnInfo) balancer.Picker {
	var subConns []*weightedSubConn

	for sc, info := range readySCs {
		// 从 resolver.Address 的 Attributes 中获取权重
		// 权重是在服务发现阶段写入的
		weight := getWeightFromAddress(info.Address)
		if weight < 1 {
			weight = 1 // 默认权重为 1
		}

		subConns = append(subConns, &weightedSubConn{
			subConn:         sc,
			weight:          weight,
			effectiveWeight: weight,
			currentWeight:   0,
			addr:            info.Address.Addr,
		})
	}

	return &smoothWeightedPicker{subConns: subConns}
}

// getWeightFromAddress 从 resolver.Address 中读取权重
// 权重数据是在服务发现阶段写入 Attributes 的
func getWeightFromAddress(addr resolver.Address) int {
	// 这里简化了从 Attributes 获取权重的过程
	// 实际实现中,权重是通过 addr.Attributes 获取的
	// 通常在服务发现的 Resolver 中将权重写入
	return 1 // 默认返回 1
}

// 注册负载均衡器
const SmoothWeightedBalancerName = "smooth_weighted_rr"

func RegisterSmoothWeightedBalancer() {
	base.NewBalancerBuilder(SmoothWeightedBalancerName, &smoothWeightedBuilder{}, base.Config{})
}

4.5 如何获取权重

之前在服务注册与发现章节提到过,ServiceInstance 本身是可以不断增加数据的。这一次我们加上权重数据,而后服务发现将权重数据注入到 Attributes 里面。

package registry

// ============================================================
// 服务实例增加权重字段
// ============================================================

// ServiceInstance 增加了 Weight 字段
type ServiceInstance struct {
	ServiceName string            `json:"serviceName"`
	InstanceID  string            `json:"instanceId"`
	Address     string            `json:"address"`
	Port        int               `json:"port"`
	Metadata    map[string]string `json:"metadata"`
	Weight      int               `json:"weight"` // 新增:权重字段
}

// 在服务发现阶段,将权重写入 resolver.Address 的 Attributes
// 这样负载均衡器就可以从 Attributes 中读取权重了

// 示例:在 etcd Resolver 中写入权重
/*
func (r *etcdResolver) updateClientConn(instances []*registry.ServiceInstance) {
	var addrs []resolver.Address
	for _, inst := range instances {
		addr := resolver.Address{
			Addr: fmt.Sprintf("%s:%d", inst.Address, inst.Port),
			// 将权重写入 Attributes
			Attributes: attributes.New("weight", inst.Weight),
		}
		addrs = append(addrs, addr)
	}
	r.cc.UpdateState(resolver.State{
		Addresses: addrs,
	})
}
*/

一些微服务框架设计得好,会允许用户自定义类似的数据。比如 Kratos 框架就允许在服务实例的 Metadata 中写入任意键值对,框架会自动传递给负载均衡器。


五、路由策略

5.1 负载均衡就够了吗?

回顾一下,我们是为了挑选出最合适的一个实例,然后将请求发送过去。负载均衡挑出来了负载最轻的节点,这就够了吗?

并不够,因为负载均衡并没有考虑业务需求

  • A/B 测试中,A 请求只能发到 A 节点上
  • 全链路压测中,压测流量只能发到测试节点上
  • VIP 服务中,VIP 的请求要发到更加高端的机器上
  • 联调或 DEBUG的时候,请求只能发送到一个特定的机器上

这些需求,一般统称为路由策略

5.2 常见路由策略

路由策略说明使用场景
标签路由给服务端实例打上不同的标签,将特定请求路由到具备某个标签的实例上A/B 测试、灰度发布
健康路由实时将服务实例分成健康和不健康两大类,每次发送请求只发送给标记为健康的实例故障隔离
转发路由用户直接指定来自某个客户端的请求转发到特定的服务实例上联调、DEBUG

5.3 路由策略设计与实现核心

不同的路由策略都要解决以下四个问题:

flowchart LR
    A[1.用户设置路由策略
提供转发规则] --> B[2.框架判断请求
是否命中路由策略] B --> C[3.框架筛选符合条件
的服务端实例] C --> D[4.执行负载均衡
找到目标节点并发送请求]

5.4 路由策略与负载均衡的关系

路由策略可以被看做是在负载均衡之前的一个步骤

能不能复用负载均衡接口? 在 gRPC 里面,答案是可能可以。需要注意:如果我们的实现在查找对应的路由策略时,所依赖的数据完全来自于 PickInfo,那么就可以(本质上主要依赖于 Ctx 字段)。如果我们的实现依赖于具体的请求参数(如 UserID),那么就不可以。

我们可以用一种非常简单的策略:过滤节点。即不管是什么路由,本质上都是为了在负载均衡之前,提前过滤一些节点:

  • 分组路由:先找出特定组的节点
  • 直接路由:使用特定 IP 和端口的节点
  • 健康优先路由:过滤出"健康"的节点

5.5 代码演示:实现分组功能(A/B 测试)

分组可以看做是一种特殊形态的路由策略——服务端实例主动给自己打上一个标签。我们用分组功能来实现一个简单的支持 A/B 流量分发的功能。

思路

  1. 在整个链路里面带上一个 A/B 标记
  2. 进程内整个 A/B 标记位放在 context.Context
  3. 利用负载均衡接口,在负载均衡之前,先根据 A/B 标记筛选节点
  4. 发送请求到具体的服务端节点上
package loadbalance

import (
	"context"
	"sync"

	"google.golang.org/grpc/balancer"
	"google.golang.org/grpc/balancer/base"
	"google.golang.org/grpc/metadata"
)

// ============================================================
// 分组路由实现:支持 A/B 流量分发
// ============================================================

// 分组相关的 context key
type groupKey struct{}

// WithGroup 将分组信息写入 context
// 调用方使用此函数设置当前请求的分组
// 例如:ctx = WithGroup(ctx, "A") 表示此请求属于 A 组
func WithGroup(ctx context.Context, group string) context.Context {
	return context.WithValue(ctx, groupKey{}, group)
}

// GetGroup 从 context 中读取分组信息
// 负载均衡器使用此函数获取当前请求应该发到哪个组
func GetGroup(ctx context.Context) string {
	if v, ok := ctx.Value(groupKey{}).(string); ok {
		return v
	}
	return "" // 没有设置分组,返回空字符串
}

// groupedSubConn 带分组信息的子连接
type groupedSubConn struct {
	subConn balancer.SubConn
	group   string // 该实例所属的分组,如 "A" 或 "B"
}

// groupPicker 分组路由 + 负载均衡 Picker
// 先按分组过滤,再在组内做轮询
type groupPicker struct {
	mu              sync.Mutex
	groupedSubConns map[string][]balancer.SubConn // 按分组组织的连接列表
	index           uint64                         // 轮询计数器
}

// Pick 实现 Picker 接口
func (p *groupPicker) Pick(info balancer.PickInfo) (balancer.PickResult, error) {
	// 第一步:从 context 中获取请求的分组标记
	group := GetGroup(info.Ctx)

	p.mu.Lock()
	defer p.mu.Unlock()

	// 第二步:根据分组筛选可用的连接
	// 如果指定了分组,只从该分组的连接中选择
	// 如果没有指定分组,从所有连接中选择
	var candidates []balancer.SubConn
	if group != "" {
		candidates = p.groupedSubConns[group]
	} else {
		// 没有指定分组,使用所有连接
		for _, conns := range p.groupedSubConns {
			candidates = append(candidates, conns...)
		}
	}

	if len(candidates) == 0 {
		return balancer.PickResult{}, balancer.ErrNoSubConnAvailable
	}

	// 第三步:在候选连接中做轮询
	idx := p.index % uint64(len(candidates))
	p.index++

	return balancer.PickResult{
		SubConn: candidates[idx],
	}, nil
}

// groupPickerBuilder 分组路由 Builder
type groupPickerBuilder struct{}

// Build 实现 PickerBuilder 接口
// 从 SubConnInfo 中提取分组信息
func (b *groupPickerBuilder) Build(readySCs map[balancer.SubConn]balancer.SubConnInfo) balancer.Picker {
	// 按分组组织连接
	groupedConns := make(map[string][]balancer.SubConn)
	for sc, info := range readySCs {
		// 从 Address 的 Attributes 中获取分组信息
		// 分组信息是在服务注册时写入的
		group := getGroupFromAddress(info.Address)
		groupedConns[group] = append(groupedConns[group], sc)
	}

	return &groupPicker{
		groupedSubConns: groupedConns,
	}
}

// getGroupFromAddress 从 resolver.Address 中读取分组信息
// 分组信息是在服务注册阶段写入的
func getGroupFromAddress(addr resolver.Address) string {
	// 实际实现中从 addr.Attributes 获取
	// 这里简化处理
	return ""
}

// ============================================================
// 服务端注册时设置分组
// ============================================================

/*
服务端在启动时通过环境变量或配置文件指定分组

```go
// 通过环境变量设置分组
group := os.Getenv("SERVICE_GROUP") // "A" 或 "B"

// 注册到 etcd 时,将分组信息写入 Metadata
instance := &registry.ServiceInstance{
	ServiceName: "user-service",
	Address:     "172.10.0.1",
	Port:        9091,
	Metadata: map[string]string{
		"group": group, // 分组信息
	},
}

*/

// ============================================================ // 客户端调用时指定分组 // ============================================================

/* 客户端在发起 RPC 调用时,通过 context 传递分组标记:

// 设置分组为 A,此请求只会发送到 A 组的实例
ctx := WithGroup(context.Background(), "A")
resp, err := client.GetUser(ctx, &GetUserRequest{Id: 123})

// 设置分组为 B,此请求只会发送到 B 组的实例
ctx = WithGroup(context.Background(), "B")
resp, err = client.GetUser(ctx, &GetUserRequest{Id: 456})

*/


> 这个图很类似于之前讨论的链路超时控制。因为本质上两者都是在整条链路里面传递一些元数据,然后根据元数据来执行一些动作。

### 5.6 过滤功能对负载均衡的影响

首先从实现的角度来说,**大部分负载均衡算法都受到了影响**,包括随机、轮询,以及对应的加权版本。

过滤功能使用不当可能会造成负载均衡算法失效:

- **过滤条件太苛刻**,以至于满足条件的实例几乎没有
- **每次请求过滤之后的节点都不同**,那么可能导致所有的请求都发到了少部分实例上

---

## 六、集群抽象(Cluster)

### 6.1 Cluster 概念

**Cluster** 是沿用 Dubbo 中的说法。它是指我们在调用远程服务的时候,尝试解决:

| 模式 | 说明 | 使用场景 |
|------|------|----------|
| **failover** | 引入重试功能,但重试时会换一个新节点 | 读操作、幂等操作 |
| **failfast** | 立刻失败,不重试 | 非幂等操作(如扣款) |
| **广播(Broadcast)** | 将请求发送到所有节点上 | 缓存刷新、配置更新 |
| **组播(Multicast)** | 将请求发送到一组节点上 | 分组通知 |

> 实际上,广播和组播可以看做是一类职能,failover 和 failfast 是另外一种职能。
>
> 注意:这里的组播和分组是不一样的——分组依旧是只发过去一个节点,而组播是发给一组节点。

### 6.2 Dubbo-go 的 Cluster 设计

在 Dubbo-go 里面,它认为 Cluster 是发起调用的一个环节,可以看做是洋葱模式里面的一层洋葱。

Cluster 相关的接口主要有两个:

- **Invoker 接口**:最核心的接口,不同 Cluster 有不同的 Invoker 实现
- **Cluster 接口**:可以看做是一层皮,将 Invoker 接口进行了封装,将多个 Cluster 组合在一起

### 6.3 failover(故障转移)

failover 的要点在于:

1. **重试**:要考虑控制重试次数和重试间隔
2. **去除失败节点,选用新节点**:
   - 新节点可能只是随便挑一个
   - 新节点也可能是不同机房上的节点
   - 新节点也可能在不同的城市

> 大多数 failover 的实现并没有那么精致,无非就是所有节点里面,去除已经使用但失败了的节点,剩下的随便挑一个。
>
> 而真的容错的话,是要考虑换机房换城市。例如如果我们知道当下调用失败是因为上海机房的网络已经崩溃了,那么我们就不应该选用任何上海机房节点,而是选用另外一个机房上的节点。这部分一般在流量调度和多活设计里面也会涉及到。

#### failover 在 gRPC 中的实现

gRPC 天然支持了重试,我们只需要提供配置。而且 gRPC 每一次重试,都是要再一次经过负载均衡的,所以实际上我们只需要:

1. 设置重试策略
2. 选用合适的负载均衡算法(如轮询)

```go
package cluster

import (
	"google.golang.org/grpc"
	"google.golang.org/grpc/backoff"
)

// ============================================================
// gRPC failover 配置示例
// ============================================================

// FailoverServiceConfig 是 gRPC 的 failover 服务配置
// 通过 ServiceConfig 设置重试策略
const FailoverServiceConfig = `{
	"methodConfig": [{
		"name": [{"service": "user.UserService"}],
		"retryPolicy": {
			"maxAttempts": 3,            // 最大重试次数(含首次调用),即最多重试 2 次
			"initialBackoff": "0.1s",    // 初始退避时间
			"maxBackoff": "1s",          // 最大退避时间
			"backoffMultiplier": 2.0,    // 退避倍数(每次退避时间乘以这个值)
			"retryableStatusCodes": [    // 哪些状态码会触发重试
				"UNAVAILABLE",
				"DEADLINE_EXCEEDED"
			]
		}
	}]
}`

// NewFailoverClient 创建支持 failover 的 gRPC 客户端
func NewFailoverClient(target string) (*grpc.ClientConn, error) {
	conn, err := grpc.Dial(
		target,
		// 指定使用轮询负载均衡
		// 因为重试时需要换节点,所以不能用哈希(哈希会选同一个节点)
		grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
		// 设置连接退避策略
		grpc.WithConnectParams(grpc.ConnectParams{
			Backoff: backoff.Config{
				BaseDelay:  100 * 1000000, // 100ms
				Multiplier: 1.6,            // 退避倍数
				Jitter:     0.2,            // 抖动因子,避免重试风暴
				MaxDelay:   3 * 1000000000, // 3s
			},
		}),
	)
	return conn, err
}

注意:在 gRPC 的负载均衡接口里,我们无法判断请求是不是重试的请求,所以我们只能选择那些每次都选出不同节点的负载均衡算法。而如果是哈希之类的负载均衡算法,同一个请求选中的都是同一个节点,所以就没办法达成 failover 的效果。

6.4 failfast(快速失败)

failfast 就是不重试——调用失败就立刻返回错误。适用于非幂等操作(如扣款、下单),重试可能导致重复操作。

// failfast 实际上就是默认行为,不设置 retryPolicy 即可
// 也可以在调用时通过 grpc.FailFast(true) 明确指定(这是默认行为)

// const FailfastServiceConfig = `{}`
// 不设置 retryPolicy 就是 failfast

6.5 广播(Broadcast)

广播在 RPC 里面是指将请求发到所有节点上。它在 gRPC 里面比较难实现,核心在于 gRPC 暴露的接口并不合适。

在 gRPC 里面,我们能够尝试的接口就是拦截器 Interceptor,它对标我们在 Web 和 ORM 里面学的 AOP 方案。

flowchart TB
    C[客户端发起调用] --> I[拦截器拦截请求]
    I --> R[从注册中心获取所有实例]
    R --> L1[调用实例1]
    R --> L2[调用实例2]
    R --> L3[调用实例3]
    L1 --> AG[聚合响应]
    L2 --> AG
    L3 --> AG
    AG --> CR[返回结果给客户端]

gRPC Interceptor

gRPC 的 Interceptor 分成好几种:

类型说明
UnaryClientInterceptor拦截 gRPC unary(一元)请求
StreamClientInterceptor拦截 gRPC stream 请求

可以看到,invoker 就接近于我们自己设计里面的 next。而为了实现广播的目标,关键就在于这个 cc *grpc.ClientConn

ClientConn 的局限

ClientConn 是 gRPC 里面的一个核心结构,可以将它理解为对某个服务的连接。因为服务本身会有很多节点,而且每个节点又可以有多个连接,所以可以将它看成是一个两层结构。

按照我们的广播目标,我们希望能够遍历三个服务实例——也就是希望能够从 ClientConn 里面拿出来对应不同实例的连接,然后发起调用。

很可惜我们拿不到。 gRPC 没有暴露从 ClientConn 中获取具体子连接的接口。我们只能考虑另辟蹊径。

gRPC 广播实现:注册中心获取所有节点

思路:

  1. 利用拦截器捕获调用
  2. 利用注册中心来获得所有的服务端实例
  3. 在拦截器内遍历所有的服务端实例,分别发起调用
package cluster

import (
	"context"
	"reflect"
	"sync"

	"google.golang.org/grpc"
)

// ============================================================
// gRPC 广播拦截器实现
// ============================================================

// broadcastKey 用于在 context 中标记广播请求
type broadcastKey struct{}

// WithBroadcast 在 context 中设置广播标记
// 只有设置了此标记的请求才会执行广播逻辑
func WithBroadcast(ctx context.Context) context.Context {
	return context.WithValue(ctx, broadcastKey{}, true)
}

// isBroadcast 检查是否为广播请求
func isBroadcast(ctx context.Context) bool {
	v, ok := ctx.Value(broadcastKey{}).(bool)
	return ok && v
}

// BroadcastInterceptor 广播拦截器
// 参数 registry 是注册中心实例,用于获取所有服务端节点
// 参数 newClient 是创建到指定实例的 gRPC 客户端的函数
func BroadcastInterceptor(
	registry Registry,        // 注册中心,用于获取所有实例
	newClient func(string) (*grpc.ClientConn, error), // 根据地址创建连接
) grpc.UnaryClientInterceptor {
	return func(
		ctx context.Context,       // 请求上下文
		method string,             // 调用的方法名,如 "/user.UserService/GetUser"
		req, reply interface{},    // 请求和响应
		cc *grpc.ClientConn,       // 客户端连接(广播时不使用)
		invoker grpc.UnaryInvoker, // 原始调用函数
		opts ...grpc.CallOption,   // 调用选项
	) error {
		// 第一步:检查是否为广播请求
		if !isBroadcast(ctx) {
			// 不是广播请求,正常调用
			return invoker(ctx, method, req, reply, cc, opts...)
		}

		// 第二步:从注册中心获取所有服务实例
		// target 是当前连接的服务名,需要从 cc 中解析
		serviceName := getServiceName(cc)
		instances, err := registry.ListServices(ctx, serviceName)
		if err != nil {
			return err
		}

		// 第三步:并发调用所有实例
		var wg sync.WaitGroup
		errs := make([]error, len(instances))

		for i, inst := range instances {
			wg.Add(1)
			go func(idx int, addr string) {
				defer wg.Done()

				// 为每个实例创建独立的连接
				// 注意:这里的连接难以复用,所以广播性能不好
				conn, err := newClient(addr)
				if err != nil {
					errs[idx] = err
					return
				}
				defer conn.Close()

				// 为每个实例创建独立的响应对象
				// 使用反射创建新对象,避免覆盖问题
				individualReply := reflect.New(reflect.ValueOf(reply).Elem().Type()).Interface()

				// 发起调用
				err = invoker(ctx, method, req, individualReply, conn, opts...)
				errs[idx] = err
			}(i, inst.Address+":"+intToStr(inst.Port))
		}

		wg.Wait()

		// 第四步:处理响应
		// 广播模式下,通常只返回第一个成功的响应
		// 也可以选择返回所有响应或丢弃所有响应
		for _, e := range errs {
			if e == nil {
				return nil // 至少一个成功就返回成功
			}
		}

		// 所有实例都失败了,返回最后一个错误
		if len(errs) > 0 {
			return errs[len(errs)-1]
		}
		return nil
	}
}

// ============================================================
// 广播响应处理策略
// ============================================================

// BroadcastResponseStrategy 广播响应处理策略
type BroadcastResponseStrategy int

const (
	// StrategyDiscard 直接丢弃响应(适用于通知类广播,如缓存刷新)
	StrategyDiscard BroadcastResponseStrategy = iota
	// StrategyFirst 返回最先收到的响应(适用于寻找最快节点)
	StrategyFirst
	// StrategyAll 返回所有响应(适用于需要汇总结果的场景)
	StrategyAll
)

// ============================================================
// 使用 channel 返回所有响应的实现
// ============================================================

// BroadcastWithChannel 使用 channel 返回所有广播响应
// 调用方可以从 channel 中读取每个实例的响应
func BroadcastWithChannel(
	ctx context.Context,
	registry Registry,
	method string,
	req interface{},
	newClient func(string) (*grpc.ClientConn, error),
	invoker grpc.UnaryInvoker,
) (<-chan interface{}, error) {
	// 获取所有实例
	serviceName := ctx.Value("serviceName").(string)
	instances, err := registry.ListServices(ctx, serviceName)
	if err != nil {
		return nil, err
	}

	// 创建响应 channel,缓冲大小等于实例数量
	ch := make(chan interface{}, len(instances))

	go func() {
		defer close(ch) // 重要:关闭 channel,让接收方知道没有更多数据了

		var wg sync.WaitGroup
		for _, inst := range instances {
			wg.Add(1)
			go func(addr string) {
				defer wg.Done()

				conn, err := newClient(addr)
				if err != nil {
					return
				}
				defer conn.Close()

				// 使用反射创建新的响应对象
				reply := reflect.New(reflect.TypeOf(req).Elem()).Interface()
				_ = invoker(ctx, method, req, reply, conn)

				// 将响应写入 channel
				// 注意:这里不会阻塞,因为 channel 有足够的缓冲
				ch <- reply
			}(inst.Address)
		}

		wg.Wait()
	}()

	return ch, nil
}

// 辅助函数
func getServiceName(cc *grpc.ClientConn) string {
	// 从 ClientConn 中解析服务名
	// 实际实现取决于具体的服务发现方案
	return ""
}

func intToStr(n int) string {
	if n == 0 {
		return "0"
	}
	var buf []byte
	if n < 0 {
		buf = append(buf, '-')
		n = -n
	}
	var digits []byte
	for n > 0 {
		digits = append([]byte{byte('0' + n%10)}, digits...)
		n /= 10
	}
	buf = append(buf, digits...)
	return string(buf)
}

广播响应处理有三种策略:

  • 直接丢弃:接近默认实现,适用于通知类广播(如缓存刷新)
  • 返回最先的:所谓的最快调用,第一个返回的响应会被取走,剩余的直接丢弃
  • 返回所有的:使用 channel 或切片传递全部响应

注意,使用 channel 传递时要防止 goroutine 泄露,同时要注意关闭 channel,否则用户在接收的时候,会不知道还有没有数据。

6.6 集群、路由和负载均衡的关系

本质上,它们都在回答同一个问题:我要把请求发给谁?

flowchart LR
    R[路由策略
筛选符合条件的节点] --> L[负载均衡
从候选节点中选择一个] L --> C[Cluster
决定调用策略] C --> S[服务端实例]
  • 路由是在负载均衡之前的一个步骤——先筛选出符合条件的节点
  • 负载均衡是从候选节点中选择一个最合适的
  • Cluster是决定调用策略——是只调一个节点(failover/failfast),还是调多个节点(广播/组播)

一般来说,Cluster 是不需要考虑负载均衡的——无论是组播还是广播,都是发给多个节点。但是 Cluster 中的组播,也可以理解为广播 + 路由,因为路由的本质就是筛选出节点。


七、总结与面试要点

7.1 负载均衡算法总结

算法是否考虑处理能力负载指标假设适用场景
轮询所有实例能力相同实例配置相同
加权轮询是(权重)权重代表处理能力实例配置不同
随机概率均匀分布简单场景
加权随机是(权重)权重代表概率实例配置不同
哈希哈希均匀需要会话保持
一致性哈希哈希均匀需要会话保持 + 动态扩缩容
最小连接数连接数连接数代表负载连接不复用的场景
最少活跃数请求数活跃数代表负载可统计活跃数的场景
最快响应响应时间响应时间代表负载对延迟敏感的场景

7.2 权重的效果总结

  • 大多数时候我们都是使用权重来表达服务器的处理能力,或者说重要性
  • 使用权重的算法都要考虑:
    • 某个实例的权重特别大,可能连续几次都选中它,要考虑平滑效果
    • 结合实际调用结果来动态调整权重,例如实例返回错误则降低权重,反之增加权重
    • 权重动态调整时要考虑上限和下限,防止权重变成 0(永不选中)或过大(永远选中)

7.3 微服务框架的局限性

在缺乏全局信息的情况下,客户端会选择服务端 1 作为服务提供者。在微服务中选择负载均衡算法,这种需要全局信息的算法可能抖动会比较厉害

那么为什么它们运作得还是很好呢?因为请求数量多了,慢慢会收敛到一种比较均匀的状态。

客户端负载均衡 vs 网关负载均衡:最大的优点是网关可以具备全局信息。如果所有的客户端都经过网关才能和服务端进行通信,那么网关就可以考虑采集所有服务端实例的负载信息,做到一个全局最优的流量调度。然而现在大多数的网关并没有利用自己的这个优势。

7.4 服务发现与筛选节点

理论上来说,在发起调用之后,调用结果要反馈给服务发现组件、Cluster 组件、路由组件和负载均衡组件,以确保:

  • 如果调用失败,这些组件要考虑将目标节点挪出可用列表(一般由服务发现负责)
  • 如果服务发现组件发现某个节点临时不可用,过一段时间之后要重新尝试探查这个节点是否恢复

为什么说值得吹嘘? 因为实际上,绝大多数微服务框架并没有做这种反馈式的服务发现和节点筛选功能。这意味着,例如在随机负载均衡里面,如果上一个请求失败了,下一个请求还是可能发给同一个节点。

7.5 面试要点

集群(Cluster)方面

  • 服务端调不通怎么办? 仔细讨论 fail-fast 还是 fail-over,以及 fail-over 要不要换节点
  • 什么是 fail-over? 就是重试,深入讨论重试次数、退避算法、换不换节点等问题
  • 什么是微服务广播? 将请求发到所有节点上。使用场景:缓存刷新(通知)、寻找最快节点(尽快服务)
  • 什么是微服务组播? 发送请求到一部分节点上,怎么分组是纯粹的业务问题

路由和负载均衡方面

  • 掌握所有主流的负载均衡算法,注意分析优缺点。分析优缺点时要注意不同算法的假设,这些假设直接决定了算法的缺点
  • 怎么设置服务器权重? 根据服务器的处理能力来设置。注意讨论在负载均衡算法里面怎么动态调整权重,同时要注意权重调整不要超过一定的边界(溢出问题)
  • 为什么有了负载均衡还是会出现某台机器被打爆? 因为那些假设实际上并不成立,例如大商家的请求就是要消耗更多资源
  • 客户端负载均衡和网关负载均衡有什么区别? 网关负载均衡有全局信息,客户端负载均衡没有,这是本质区别
  • 什么是微服务路由? 根据用户设置,将符合条件的请求发送到特定节点的过程。强调分组可以看做是一种特殊的路由策略
  • 如何自定义负载均衡算法? 根据业务特征随便挑几个指标来设计——错误率、CPU、IO、网络负载等

设计自己的负载均衡算法

核心是根据自己的业务特征来选取一些指标,来表达服务实例的负载:

  • 错误率等服务指标
  • CPU、IO、网络负载等硬件指标

而想要知道这些指标,除了客户端统计,还有一些奇技淫巧:

  • 服务端将指标的值写入注册中心,注册中心通知客户端
  • 服务端每次返回响应的时候,额外带上自己的指标(如 CPU 利用率)
  • 利用可观测性平台,从观测性平台获得数据

总结

本教程从负载均衡的基本概念出发,讲解了各类负载均衡算法的原理和实现,然后在 gRPC 框架中实现了自定义负载均衡器,接着介绍了路由策略和分组功能的实现,最后讨论了集群抽象中的 failover、failfast、广播和组播。

关键知识点回顾:

  1. 负载均衡算法分为不实时计算(轮询、随机、哈希等)和实时计算(最小连接数、最少活跃数、最快响应)两大类
  2. 每种算法都有假设,理解假设是理解算法优缺点的关键
  3. 平滑加权轮询通过动态调整 currentWeight 来避免权重大的节点被连续选中
  4. 一致性哈希通过环状结构和虚拟节点解决了扩缩容时大量请求重新分配的问题
  5. gRPC 负载均衡需要实现 PickerBuilder 和 Picker 两个接口
  6. 路由策略是负载均衡之前的过滤步骤,解决业务需求(A/B 测试、灰度发布等)
  7. Cluster 模式包括 failover(重试换节点)、failfast(快速失败)、广播(发给所有节点)、组播(发给一组节点)
  8. 广播在 gRPC 中难以直接实现,需要通过拦截器 + 注册中心的变通方案
  9. 客户端负载均衡缺乏全局信息,但请求数多了会收敛到均匀状态

下一章我们将讲解可用性和可观测性,探讨如何保证微服务的稳定运行和问题排查。


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 负载均衡算法分"不实时计算"和"实时计算"两大类。请各举两个代表算法,并说明这两类在假设上的根本区别。
  2. 为什么权重很大的节点在简单加权轮询下会被"连续选中 5 次",而平滑加权轮询能避免?关键改动是哪一步?
  3. 普通哈希在实例扩缩容时为什么会导致"大部分请求重新分配"?一致性哈希靠什么机制把影响范围缩小到"一部分请求"?
  4. 在 gRPC 里做 failover(重试换节点)时,为什么不能使用哈希类负载均衡算法?该选哪种?
  5. 路由策略与负载均衡是什么关系?为什么说"路由本质上是负载均衡之前的一个过滤步骤"?

动手练习(建议真做一遍)

  1. 把文中 RoundRobinPickerSmoothWeightedRRPicker 拷到一个小 main 里,权重设为 [5,1,1],各 Pick 12 次,比较两次输出序列的平滑程度。
  2. 在一致性哈希环上手动放 4 个节点(哈希值 1000/3000/7000/9000),分别测试 key 哈希为 1500、5000、8500 时落在哪个节点;然后"新增一个节点 5000",重新计算,观察有多少 key 的归属发生了变化。
  3. 仿照 5.5 的分组路由,给 groupPicker 增加"健康优先"过滤:在 groupedSubConns 之外维护一份不健康节点集合,Pick 时优先排除它们,验证"部分节点不健康时流量自动绕开"。

本章小结

  • 两类算法,两类假设:不实时计算的算法(轮询/随机/哈希等)依赖"统计规律会收敛",实时计算的算法(最少活跃数/最快响应)依赖"能拿到准确的实时负载"——但实时数据本身有开销且可能过时,所以大多数场景仍用前者。
  • 权重是表达能力的手段:加权算法要关注平滑效果与动态上下限,防止权重变成 0(永不选中)或过大(永远选中)。
  • 一致性哈希解决扩缩容痛点:环状结构 + 虚拟节点让增删节点只影响部分 key,是"会话保持 + 动态扩缩容"场景的优选。
  • gRPC 自定义负载均衡 = 两个接口PickerBuilder 构建 PickerPicker.Pick 每次 RPC 选连接;通过 ServiceConfig 指定。
  • 路由在负载均衡之前:路由按业务规则(标签/健康/分组)先过滤节点,负载均衡再从候选里选一个;Cluster 再决定调一个还是多个节点(failover/failfast/广播/组播)。
  • 过渡到下一章:下一章讲解可用性与可观测性,看看如何保证微服务稳定运行并快速排查问题。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!