Spring Cloud LoadBalancer 源码阅读
Spring Cloud LoadBalancer 是 Ribbon 的官方接班人,核心是一个响应式接口 ReactorLoadBalancer。本文从源码拆解:负载均衡算法如何实现、Reactive 与阻塞式如何桥接、实例列表缓存如何构建与更新。
核心接口体系
LoadBalancer<T>(顶层接口)
└── ReactorLoadBalancer<T>
├── RoundRobinLoadBalancer 轮询
├── RandomLoadBalancer 随机
├── WeightedLoadBalancer? 权重(由 supplier 组合实现)
└── NacosLoadBalancer(Alibaba 扩展)源码接口
java
// org.springframework.cloud.loadbalancer.core.ReactorLoadBalancer
public interface ReactorLoadBalancer<T> extends LoadBalancer<T> {
Mono<Response<T>> choose(Request request);
}
// 默认请求实现
public class DefaultRequest implements Request {
private final RequestDataContext context; // 请求上下文
private final String serviceId; // 目标服务名
}choose 返回 Mono<Response<T>>,Response 包含选中的实例或错误状态。
轮询算法:RoundRobinLoadBalancer
java
// org.springframework.cloud.loadbalancer.core.RoundRobinLoadBalancer
public class RoundRobinLoadBalancer implements ReactorServiceInstanceLoadBalancer {
private final AtomicInteger position; // 原子计数器,记录当前轮询位置
public Mono<Response<ServiceInstance>> choose(Request request) {
ServiceInstanceListSupplier supplier = serviceInstanceListSupplierProvider
.getIfAvailable(NoopServiceInstanceListSupplier::new);
return supplier.get().next()
.map(instances -> processInstanceResponse(supplier, instances));
}
private Response<ServiceInstance> processInstanceResponse(
ServiceInstanceListSupplier supplier, List<ServiceInstance> instances) {
Response<ServiceInstance> serviceInstanceResponse = getInstanceResponse(instances);
if (supplier instanceof HealthCheckServiceInstanceListSupplier) {
// 健康检查 supplier 会处理过滤后的结果
}
return serviceInstanceResponse;
}
private Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
if (instances.isEmpty()) return emptyResponse();
int pos = Math.abs(this.position.incrementAndGet()); // 计数器自增
int index = pos % instances.size(); // 取模选下标
return new DefaultResponse(instances.get(index)); // 轮询核心:取模
}
}轮询的核心就一句话:计数器自增后对实例数取模。AtomicInteger 保证多线程下的原子自增。
空白响应
java
private Response<ServiceInstance> emptyResponse() {
return new EmptyResponse(); // 实例为空时的响应,触发重试或报错
}随机算法:RandomLoadBalancer
java
// org.springframework.cloud.loadbalancer.core.RandomLoadBalancer
public class RandomLoadBalancer implements ReactorServiceInstanceLoadBalancer {
private final AtomicInteger position; // 仅用于配合缓存刷新判断
private Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
if (instances.isEmpty()) return emptyResponse();
int index = ThreadLocalRandom.current().nextInt(instances.size()); // 随机下标
return new DefaultResponse(instances.get(index));
}
}随机算法用 ThreadLocalRandom 生成下标,无共享状态竞争。
权重算法
Spring Cloud LoadBalancer 没有独立的 WeightedLoadBalancer,权重通过 supplier 组合实现:
WeightedServiceInstanceListSupplier
java
// org.springframework.cloud.loadbalancer.core.WeightedServiceInstanceListSupplier
public class WeightedServiceInstanceListSupplier
extends DelegatingServiceInstanceListSupplier implements ServiceInstanceListSupplier {
public Flux<List<ServiceInstance>> get() {
return delegate.get().map(this::weightedInstances);
}
private List<ServiceInstance> weightedInstances(List<ServiceInstance> instances) {
if (instances.isEmpty()) return instances;
List<WeightedServiceInstance> weighted = instances.stream()
.map(WeightedServiceInstance::new) // 解析 metadata 中的 weight
.collect(toList());
// 按权重累加区间,随机落在哪个区间就选哪个实例
WeightedServiceInstance winner = ...;
return Collections.singletonList(winner);
}
}实现逻辑:
实例 A weight=1 [0, 1)
实例 B weight=2 [1, 3)
实例 C weight=3 [3, 6) ← 随机数落在 [0,6),按区间命中权重取自实例 metadata 中的 weight 字段(注册时写入)。
Reactive 与阻塞式桥接
阻塞式接口:BlockingLoadBalancerClient
java
// org.springframework.cloud.client.loadbalancer.BlockingLoadBalancerClient
public class BlockingLoadBalancerClient implements LoadBalancerClient {
private final LoadBalancerClientFactory loadBalancerClientFactory;
@Override
public ServiceInstance choose(String serviceId) {
// 1. 从 factory 拿到该服务的 ReactorLoadBalancer
ReactorLoadBalancer<ServiceInstance> loadBalancer =
(ReactorLoadBalancer<ServiceInstance>) loadBalancerClientFactory
.getInstance(serviceId, ReactorServiceInstanceLoadBalancer.class);
// 2. 同步阻塞等待 choose 结果
Response<ServiceInstance> loadBalancerResponse =
Mono.from(loadBalancer.choose(request)).block();
...
return loadBalancerResponse.getServer();
}
}关键点:阻塞式是 Reactive 的同步壳。choose 内部是响应式调用,阻塞版用 .block() 等待结果。OpenFeign、@LoadBalanced RestTemplate 都走这个入口。
LoadBalancerClientFactory
java
// org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory
public class LoadBalancerClientFactory extends NamedContextFactory<LoadBalancerClientSpecification> {
public LoadBalancerClientFactory() {
super(LoadBalancerClientConfiguration.class, "loadbalancer", "spring.cloud.loadbalancer");
}
}NamedContextFactory 为每个服务名创建独立子上下文:不同服务可以配置不同的负载均衡器与 supplier,互不干扰。这与 Feign 的 FeignClientFactory 是同一套机制。
服务实例缓存与更新
CachingServiceInstanceListSupplier
实例列表默认有缓存,避免每次请求都查注册中心:
java
// org.springframework.cloud.loadbalancer.core.CachingServiceInstanceListSupplier
public class CachingServiceInstanceListSupplier
extends DelegatingServiceInstanceListSupplier implements ServiceInstanceListSupplier {
private final Flux<List<ServiceInstance>> cached; // 缓存流
private final CacheManager cacheManager; // 底层缓存(如 Caffeine)
private final String serviceId;
public CachingServiceInstanceListSupplier(...) {
super(delegate);
// 设置 TTL(默认 35s),定时刷新
this.cached = cacheManager.getCache("instances")
.get(serviceId, ...)
.timeout(Duration.ofSeconds(...));
}
@Override
public Flux<List<ServiceInstance>> get() {
return cached; // 直接返回缓存
}
}刷新机制
| 机制 | 说明 |
|---|---|
| TTL 过期 | 默认 35 秒后缓存失效,重新从 DiscoveryClient 拉取 |
| 主动刷新 | -Dspring.cloud.loadbalancer.cache.ttl=10 调短刷新间隔 |
| 禁用缓存 | -Dspring.cloud.loadbalancer.cache.enabled=false(每次实时查) |
配置项:
yaml
spring:
cloud:
loadbalancer:
cache:
enabled: true
ttl: 35s # 缓存过期时间
capacity: 256 # 缓存容量更新链路
注册中心实例变化
→ DiscoveryClient(Nacos 长轮询/推送)
→ 缓存 TTL 到期重新 get()
→ ServiceInstanceListSupplier 链重新拉取
→ 轮询/随机算法基于新列表选取扩展:自定义负载均衡算法
实现接口并注册即可:
java
@Component
public class HashLoadBalancer implements ReactorServiceInstanceLoadBalancer {
@Override
public Mono<Response<ServiceInstance>> choose(Request request) {
return Mono.defer(() -> {
String serviceId = ((RequestDataContext) request.getContext()).getClientRequest().getUrl().getHost();
List<ServiceInstance> instances = ...; // 从 supplier 获取
int index = Math.abs(serviceId.hashCode()) % instances.size();
return Mono.just(new DefaultResponse(instances.get(index)));
});
}
}通过 LoadBalancerClientFactory 或 @LoadBalancerClients 配置替换默认实现。
关键源码文件索引
| 类 | 位置 |
|---|---|
| ReactorLoadBalancer | spring-cloud-loadbalancer-core/.../ReactorLoadBalancer.java |
| RoundRobinLoadBalancer | spring-cloud-loadbalancer-core/.../RoundRobinLoadBalancer.java |
| RandomLoadBalancer | spring-cloud-loadbalancer-core/.../RandomLoadBalancer.java |
| BlockingLoadBalancerClient | spring-cloud-commons/.../BlockingLoadBalancerClient.java |
| CachingServiceInstanceListSupplier | spring-cloud-loadbalancer-core/.../CachingServiceInstanceListSupplier.java |
| LoadBalancerClientFactory | spring-cloud-loadbalancer-core/.../LoadBalancerClientFactory.java |
常见问题
- 为什么默认是轮询?
LoadBalancerClientConfiguration默认装配RoundRobinLoadBalancer,无配置时即轮询。 - 实例列表多久更新一次? 默认缓存 TTL 35 秒;实例下线后最长 35 秒内仍可能被选中(结合注册中心推送可缩短)。
- 服务名不同能配不同算法吗? 能。
NamedContextFactory为每服务独立上下文,可用@LoadBalancerClient(name="order", configuration=...)指定。 - Nacos 有自己的负载均衡吗? Alibaba 提供
NacosLoadBalancer(基于 Nacos 权重),可通过配置启用。