5. Dubbo3 集群源码分析
集群
参考链接:dubbo2.6源码分析
1. 集群模块概述
在分布式系统中,为了保证高可用性,同一服务通常会部署多个实例。对于服务消费者而言,如何从这些实例中选择一个进行调用,以及如何处理调用失败的情况,是需要特别关注的问题。Dubbo的集群模块正是为解决这些问题而设计的。
集群模块的核心职责包括:
- 将多个服务提供者合并为一个虚拟的服务提供点
- 在调用失败时,根据配置的策略进行重试或其他容错处理
- 屏蔽服务提供者集群的复杂性,简化服务消费者的调用逻辑
Dubbo通过Cluster接口和Cluster Invoker的设计模式实现了这些功能。Cluster负责将多个服务提供者合并为一个Cluster Invoker,而Cluster Invoker则负责具体的调用逻辑和容错策略。
2. 集群容错架构
在分析具体代码之前,我们先了解Dubbo集群容错的整体架构。
2.1 核心组件
Dubbo集群容错涉及以下核心组件:
- Cluster: 集群接口,将多个Invoker合并成一个Invoker
- Directory: 服务目录,存储Invoker列表,可动态感知服务提供者变化
- Router: 路由器,根据规则筛选Invoker
- LoadBalance: 负载均衡器,从多个Invoker中选择一个
- Cluster Invoker: 集群Invoker,封装了集群容错逻辑
2.2 集群工作流程
集群工作流程分为两个阶段:
- 初始化阶段:服务消费者启动时,Cluster实现类为服务消费者创建Cluster Invoker实例。
- 调用阶段:服务消费者发起远程调用时,Cluster Invoker执行以下步骤:
- 调用Directory的list方法获取Invoker列表
- 调用Router的route方法进行路由,过滤不符合条件的Invoker
- 使用LoadBalance从Invoker列表中选择一个Invoker
- 调用选中Invoker的invoke方法进行远程调用
- 根据配置的容错策略处理调用结果或异常
3. 容错策略详解
Dubbo提供了多种容错策略,以适应不同的业务场景:
| 容错策略 | 配置值 | 功能描述 | 适用场景 |
|---|---|---|---|
| Failover Cluster | failover | 失败自动切换,重试其他服务器 | 适用于幂等操作,如读取操作 |
| Failfast Cluster | failfast | 快速失败,只发起一次调用 | 适用于非幂等操作,如新增记录 |
| Failsafe Cluster | failsafe | 失败安全,出现异常时忽略 | 适用于日志记录等操作 |
| Failback Cluster | failback | 失败自动恢复,定时重发 | 适用于消息通知等操作 |
| Forking Cluster | forking | 并行调用多个服务提供者 | 适用于实时性要求较高的读操作 |
| Broadcast Cluster | broadcast | 广播调用所有提供者 | 适用于通知所有提供者更新缓存等操作 |
Dubbo3新增了以下容错策略:
| 容错策略 | 配置值 | 功能描述 | 适用场景 |
|---|---|---|---|
| ZoneAware Cluster | zone-aware | 区域感知,优先调用同区域服务 | 适用于多区域部署场景 |
| Available Cluster | available | 调用第一个可用的服务提供者 | 适用于服务状态查询等非关键操作 |
4. 源码分析
4.1 Cluster接口与实现
Cluster接口定义了将多个Invoker合并为一个Invoker的方法:
1
2
3
4
public interface Cluster {
@Adaptive
<T> Invoker<T> join(Directory<T> directory) throws RpcException;
}
AbstractCluster是所有Cluster实现类的父类:
1
2
3
4
5
6
7
8
public abstract class AbstractCluster implements Cluster {
@Override
public <T> Invoker<T> join(Directory<T> directory) throws RpcException {
return doJoin(directory);
}
protected abstract <T> AbstractClusterInvoker<T> doJoin(Directory<T> directory);
}
以FailoverCluster为例,其实现非常简单:
1
2
3
4
5
6
7
8
public class FailoverCluster extends AbstractCluster {
public static final String NAME = "failover";
@Override
public <T> AbstractClusterInvoker<T> doJoin(Directory<T> directory) throws RpcException {
return new FailoverClusterInvoker<>(directory);
}
}
4.2 AbstractClusterInvoker分析
AbstractClusterInvoker是所有Cluster Invoker的父类,封装了通用的调用逻辑:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
public Result invoke(final Invocation invocation) throws RpcException {
checkWhetherDestroyed();
// 记录路由阶段的性能分析
InvocationProfilerUtils.enterDetailProfiler(invocation, () -> "Router route.");
// 获取可用的服务提供者列表
List<Invoker<T>> invokers = list(invocation);
InvocationProfilerUtils.releaseDetailProfiler(invocation);
// 检查提供者列表有效性
checkInvokers(invokers, invocation);
// 初始化负载均衡策略
LoadBalance loadbalance = initLoadBalance(invokers, invocation);
// 为异步调用附加ID
RpcUtils.attachInvocationIdIfAsync(getUrl(), invocation);
// 记录集群调用阶段的性能分析
InvocationProfilerUtils.enterDetailProfiler(
invocation, () -> "Cluster " + this.getClass().getName() + " invoke.");
try {
// 执行具体的集群调用逻辑
return doInvoke(invocation, invokers, loadbalance);
} finally {
InvocationProfilerUtils.releaseDetailProfiler(invocation);
}
}
4.3 FailoverClusterInvoker分析
FailoverClusterInvoker是默认的集群调用策略,实现了失败自动切换功能:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
public Result doInvoke(Invocation invocation, final List<Invoker<T>> invokers, LoadBalance loadbalance)
throws RpcException {
// 创建Invoker列表副本,避免并发修改
List<Invoker<T>> copyInvokers = invokers;
// 获取方法名
String methodName = RpcUtils.getMethodName(invocation);
// 计算调用次数(包含重试次数)
int len = calculateInvokeTimes(methodName);
RpcException le = null; // 记录最后一次异常
List<Invoker<T>> invoked = new ArrayList<>(copyInvokers.size()); // 已调用的Invoker列表
Set<String> providers = new HashSet<>(len); // 记录失败的Provider地址
// 重试循环
for (int i = 0; i < len; i++) {
// 非首次调用时重新获取Invoker列表
if (i > 0) {
checkWhetherDestroyed();
copyInvokers = list(invocation);
checkInvokers(copyInvokers, invocation);
}
// 使用负载均衡选择Invoker
Invoker<T> invoker = select(loadbalance, invocation, copyInvokers, invoked);
invoked.add(invoker);
RpcContext.getServiceContext().setInvokers((List) invoked);
boolean success = false;
try {
// 执行调用
Result result = invokeWithContext(invoker, invocation);
// 记录重试成功的情况
if (le != null && logger.isWarnEnabled()) {
logger.warn("/*...*/");
}
success = true;
return result;
} catch (RpcException e) {
if (e.isBiz()) { // 业务异常直接抛出
throw e;
}
le = e;
} catch (Throwable e) {
le = new RpcException(e.getMessage(), e);
} finally {
if (!success) {
providers.add(invoker.getUrl().getAddress());
}
}
}
// 所有重试均失败,抛出异常
throw new RpcException("/*...*/");
}
4.4 粘滞连接特性
Dubbo的粘滞连接(Sticky Connection)特性允许服务消费者尽可能地使用同一个服务提供者,提高本地缓存命中率:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
protected Invoker<T> select(
LoadBalance loadbalance, Invocation invocation, List<Invoker<T>> invokers, List<Invoker<T>> selected)
throws RpcException {
if (CollectionUtils.isEmpty(invokers)) {
return null;
}
String methodName = invocation == null ? StringUtils.EMPTY_STRING : RpcUtils.getMethodName(invocation);
// 检查是否启用粘滞会话(Sticky Session)特性
// 从Invoker URL中获取方法级配置,默认使用全局配置
boolean sticky = invokers.get(0).getUrl().getMethodParameter(methodName, CLUSTER_STICKY_KEY, DEFAULT_CLUSTER_STICKY);
// 处理粘滞Invoker失效场景:
// 若缓存的粘滞Invoker不在当前可用列表中,清空缓存
if (stickyInvoker != null && !invokers.contains(stickyInvoker)) {
stickyInvoker = null;
}
// 核心粘滞会话逻辑:
// 若启用粘滞且缓存了粘滞Invoker,且未在已选列表中重复选择
if (sticky && stickyInvoker != null && (selected == null || !selected.contains(stickyInvoker))) {
// 附加可用性检查(可选配置)
if (availableCheck && stickyInvoker.isAvailable()) {
return stickyInvoker; // 直接返回粘滞Invoker
}
}
// 执行实际的负载均衡选择
Invoker<T> invoker = doSelect(loadbalance, invocation, invokers, selected);
// 更新粘滞Invoker(若启用粘滞会话)
if (sticky) {
stickyInvoker = invoker;
}
return invoker;
}
select 方法的主要逻辑集中在了对粘滞连接特性的支持上。首先是获取 sticky 配置,然后再检测 invokers 列表中是否包含 stickyInvoker,如果不包含,则认为该 stickyInvoker 不可用,此时将其置空。这里的 invokers 列表可以看做是存活着的服务提供者列表,如果这个列表不包含 stickyInvoker,那自然而然的认为 stickyInvoker 挂了,所以置空。如果 stickyInvoker 存在于 invokers 列表中,此时要进行下一项检测 — 检测 selected 中是否包含 stickyInvoker。如果包含的话,说明 stickyInvoker 在此之前没有成功提供服务(但其仍然处于存活状态)。此时我们认为这个服务不可靠,不应该在重试期间内再次被调用,因此这个时候不会返回该 stickyInvoker。如果 selected 不包含 stickyInvoker,此时还需要进行可用性检测,比如检测服务提供者网络连通性等。当可用性检测通过,才可返回 stickyInvoker,否则调用 doSelect 方法选择 Invoker。如果 sticky 为 true,此时会将 doSelect 方法选出的 Invoker 赋值给 stickyInvoker。
下面我们来看一下 doSelect 方法的具体实现:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
private Invoker<T> doSelect(
LoadBalance loadbalance, Invocation invocation, List<Invoker<T>> invokers, List<Invoker<T>> selected)
throws RpcException {
// 处理边界情况:可用Invoker列表为空
if (CollectionUtils.isEmpty(invokers)) {
return null;
}
// 处理边界情况:只有一个可用Invoker
if (invokers.size() == 1) {
Invoker<T> tInvoker = invokers.get(0);
checkShouldInvalidateInvoker(tInvoker); // 检查Invoker是否应失效
return tInvoker;
}
// 核心选择逻辑:使用负载均衡策略选择Invoker
Invoker<T> invoker = loadbalance.select(invokers, getUrl(), invocation);
// 检查选中Invoker的状态:是否已选或不可用
boolean isSelected = selected != null && selected.contains(invoker); // 是否已在已选列表中
boolean isUnavailable = availableCheck && !invoker.isAvailable() && getUrl() != null; // 是否不可用
// 处理不可用Invoker:标记为失效
if (isUnavailable) {
invalidateInvoker(invoker);
}
// 触发重试选择的条件:已选或不可用
if (isSelected || isUnavailable) {
try {
// 执行重试选择逻辑
Invoker<T> rInvoker = reselect(loadbalance, invocation, invokers, selected, availableCheck);
if (rInvoker != null) {
invoker = rInvoker; // 重试选择成功,更新Invoker
} else {
// 重试选择失败,执行故障转移:选择列表中的下一个Invoker
int index = invokers.indexOf(invoker);
try {
// 使用取模运算确保索引有效,避免越界
invoker = invokers.get((index + 1) % invokers.size());
} catch (Exception e) {
logger.warn("No available invoker in the second selection of doSelect, ignore this error.");
}
}
} catch (Throwable t) {
// 重试选择异常处理:记录错误并保持当前Invoker
logger.error("Exception when doSelect in cluster, ignore this error.", t);
}
}
return invoker;
}
doSelect 主要做了两件事,第一是通过负载均衡组件选择 Invoker。第二是,如果选出来的 Invoker 不稳定,或不可用,此时需要调用 reselect 方法进行重选。若 reselect 选出来的 Invoker 为空,此时定位 invoker 在 invokers 列表中的位置 index,然后获取 index + 1 处的 invoker,这也可以看做是重选逻辑的一部分。下面我们来看一下 reselect 方法的逻辑。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
private Invoker<T> reselect(
LoadBalance loadbalance,
Invocation invocation,
List<Invoker<T>> invokers,
List<Invoker<T>> selected,
boolean availableCheck)
throws RpcException {
// 预分配重选列表,避免动态扩容开销
// 数量限制为可用Invoker数与重选配置数的较小值
List<Invoker<T>> reselectInvokers = new ArrayList<>(Math.min(invokers.size(), reselectCount));
/*
* 阶段1:筛选未选过的可用Invoker
* 策略:
* - 当重选配置数≥可用Invoker数时,尝试选择所有符合条件的Invoker
* - 否则随机选择指定数量的Invoker(避免全量遍历大型列表)
* 目标:构建一个包含未选且可用Invoker的候选列表
*/
if (reselectCount >= invokers.size()) {
// 全量筛选模式
for (Invoker<T> invoker : invokers) {
// 可用性检查
if (availableCheck && !invoker.isAvailable()) {
invalidateInvoker(invoker); // 标记不可用Invoker
continue;
}
// 去重检查:未在已选列表中
if (selected == null || !selected.contains(invoker)) {
reselectInvokers.add(invoker);
}
}
} else {
// 随机筛选模式(避免大型列表全量遍历)
for (int i = 0; i < reselectCount; i++) {
// 随机选择Invoker
Invoker<T> invoker = invokers.get(ThreadLocalRandom.current().nextInt(invokers.size()));
// 可用性检查
if (availableCheck && !invoker.isAvailable()) {
invalidateInvoker(invoker); // 标记不可用Invoker
continue;
}
// 双重去重:不在已选列表且不在重选列表中
if ((selected == null || !selected.contains(invoker)) && !reselectInvokers.contains(invoker)) {
reselectInvokers.add(invoker);
}
}
}
/*
* 阶段2:从候选列表中通过负载均衡选择Invoker
* 条件:候选列表非空时执行负载均衡
*/
if (!reselectInvokers.isEmpty()) {
return loadbalance.select(reselectInvokers, getUrl(), invocation);
}
/*
* 阶段3:候选列表为空时的回查机制
* 场景:所有未选Invoker均不可用或已被选择
* 策略:检查已选列表中是否有可用Invoker
*/
if (selected != null) {
for (Invoker<T> invoker : selected) {
// 优先选择可用且未在候选列表中的Invoker
if (invoker.isAvailable() && !reselectInvokers.contains(invoker)) {
reselectInvokers.add(invoker);
}
}
}
/*
* 阶段4:回查结果的二次负载均衡
* 条件:回查后候选列表非空时执行
*/
if (!reselectInvokers.isEmpty()) {
return loadbalance.select(reselectInvokers, getUrl(), invocation);
}
/*
* 阶段5:最终失败处理
* 场景:所有Invoker均不可用或已被尝试
* 处理:返回null,交由上层逻辑处理
*/
return null;
}
reselect 方法总结下来其实只做了两件事情,第一是查找可用的 Invoker,并将其添加到 reselectInvokers 集合中。第二,如果 reselectInvokers 不为空,则通过负载均衡组件再次进行选择。其中第一件事情又可进行细分,一开始,reselect 从 invokers 列表中查找有效可用的 Invoker,若未能找到,此时再到 selected 列表中继续查找。关于 reselect 方法就先分析到这,继续分析其他的 Cluster Invoker。
4.5 其他容错策略分析
4.5.1 FailfastClusterInvoker
FailfastClusterInvoker只会进行一次调用,失败后立即抛出异常:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
public Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
// 选择Invoker
Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
try {
// 调用Invoker
return invokeWithContext(invoker, invocation);
} catch (Throwable e) {
if (e instanceof RpcException && ((RpcException) e).isBiz()) {
throw (RpcException) e;
}
// 抛出异常
throw new RpcException("Failfast invoke providers failed...");
}
}
4.5.2 FailsafeClusterInvoker
FailsafeClusterInvoker在调用失败时不抛出异常,而是返回空结果:
1
2
3
4
5
6
7
8
9
10
11
12
public Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
try {
// 选择Invoker
Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
// 调用Invoker
return invokeWithContext(invoker, invocation);
} catch (Throwable e) {
logger.error("Failsafe ignore exception: " + e.getMessage(), e);
// 返回空结果
return AsyncRpcResult.newDefaultAsyncResult(null, null, invocation);
}
}
4.5.3 FailbackClusterInvoker
FailbackClusterInvoker在调用失败后,通过定时任务进行重试:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
public class FailbackClusterInvoker<T> extends AbstractClusterInvoker<T> {
private static final long RETRY_FAILED_PERIOD = 5 * 1000;
private final ScheduledExecutorService scheduledExecutorService;
private final ConcurrentMap<Invocation, AbstractClusterInvoker<?>> failed;
private volatile ScheduledFuture<?> retryFuture;
@Override
protected Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
try {
// 选择Invoker并调用
Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
return invoker.invoke(invocation);
} catch (Throwable e) {
// 记录失败的调用
addFailed(invocation, this);
// 返回空结果
return new RpcResult();
}
}
private void addFailed(Invocation invocation, AbstractClusterInvoker<?> router) {
// 创建定时任务进行重试
if (retryFuture == null) {
synchronized (this) {
if (retryFuture == null) {
retryFuture = scheduledExecutorService.scheduleWithFixedDelay(
this::retryFailed, RETRY_FAILED_PERIOD, RETRY_FAILED_PERIOD, TimeUnit.MILLISECONDS);
}
}
}
failed.put(invocation, router);
}
void retryFailed() {
// 遍历失败的调用并重试
for (Map.Entry<Invocation, AbstractClusterInvoker<?>> entry :
new HashMap<>(failed).entrySet()) {
Invocation invocation = entry.getKey();
Invoker<?> invoker = entry.getValue();
try {
invoker.invoke(invocation);
failed.remove(invocation);
} catch (Throwable e) {
logger.error("Failed retry to invoke method...", e);
}
}
}
}
4.5.4 ForkingClusterInvoker
ForkingClusterInvoker会并行调用多个服务提供者,只要有一个成功即返回:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
public Result doInvoke(final Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)
throws RpcException {
// 并行调用配置
final List<Invoker<T>> selected;
final int forks = getUrl().getParameter(FORKS_KEY, DEFAULT_FORKS);
final int timeout = getUrl().getParameter(TIMEOUT_KEY, DEFAULT_TIMEOUT);
// 选择要并行调用的Invoker列表
if (forks <= 0 || forks >= invokers.size()) {
selected = invokers;
} else {
selected = new ArrayList<>(forks);
while (selected.size() < forks) {
Invoker<T> invoker = select(loadbalance, invocation, invokers, selected);
if (!selected.contains(invoker)) {
selected.add(invoker);
}
}
}
RpcContext.getServiceContext().setInvokers((List) selected);
// 结果收集
final AtomicInteger count = new AtomicInteger();
final BlockingQueue<Object> ref = new LinkedBlockingQueue<>(1);
// 并行调用
selected.forEach(invoker -> {
URL consumerUrl = RpcContext.getServiceContext().getConsumerUrl();
CompletableFuture.<Object>supplyAsync(
() -> {
if (ref.size() > 0) {
return null;
}
return invokeWithContextAsync(invoker, invocation, consumerUrl);
},
executor)
.whenComplete((v, t) -> {
if (t == null) {
ref.offer(v);
} else {
int value = count.incrementAndGet();
if (value >= selected.size()) {
ref.offer(t);
}
}
});
});
// 等待结果
try {
Object ret = ref.poll(timeout, TimeUnit.MILLISECONDS);
if (ret instanceof Throwable) {
Throwable e = (Throwable) ret;
throw new RpcException("Failed to forking invoke...");
}
return (Result) ret;
} catch (InterruptedException e) {
throw new RpcException("Failed to forking invoke...");
}
}
4.5.5 BroadcastClusterInvoker
BroadcastClusterInvoker会调用所有服务提供者,任何一个出错就抛出异常:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
public Result doInvoke(final Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)
throws RpcException {
RpcContext.getServiceContext().setInvokers((List) invokers);
RpcException exception = null;
Result result = null;
// 获取广播失败阈值配置
int broadcastFailPercent = getUrl().getParameter(BROADCAST_FAIL_PERCENT_KEY, MAX_BROADCAST_FAIL_PERCENT);
// 计算失败阈值
int failThresholdIndex = invokers.size() * broadcastFailPercent / MAX_BROADCAST_FAIL_PERCENT;
int failIndex = 0;
// 调用所有服务提供者
for (int i = 0; i < invokers.size(); i++) {
Invoker<T> invoker = invokers.get(i);
RpcContext.RestoreContext restoreContext = new RpcContext.RestoreContext();
try {
// 创建子调用
RpcInvocation subInvocation = new RpcInvocation(/*...*/);
// 执行调用
result = invokeWithContext(invoker, subInvocation);
// 处理异常结果
if (result.hasException()) {
exception = getRpcException(result.getException());
failIndex++;
if (failIndex == failThresholdIndex) {
break;
}
}
} catch (Throwable e) {
exception = getRpcException(e);
failIndex++;
if (failIndex == failThresholdIndex) {
break;
}
} finally {
if (i != invokers.size() - 1) {
restoreContext.restore();
}
}
}
// 处理异常
if (exception != null) {
throw exception;
}
return result;
}
4.6 Dubbo3新增容错策略
4.6.1 ZoneAwareClusterInvoker
Dubbo3新增的ZoneAwareClusterInvoker实现了区域感知功能,优先调用同区域的服务提供者:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
public class ZoneAwareClusterInvoker<T> extends AbstractClusterInvoker<T> {
@Override
protected Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
// 获取当前区域
String zone = ZoneHelper.getZone(getUrl());
if (StringUtils.isEmpty(zone)) {
// 没有区域信息时,使用Failover策略
return failoverInvoker.invoke(invocation);
}
// 按区域分组
Map<String, List<Invoker<T>>> zoneToInvokers = new HashMap<>();
for (Invoker<T> invoker : invokers) {
String invokerZone = ZoneHelper.getZone(invoker.getUrl());
zoneToInvokers.computeIfAbsent(invokerZone, k -> new ArrayList<>()).add(invoker);
}
// 优先使用同区域Invoker
List<Invoker<T>> sameZoneInvokers = zoneToInvokers.get(zone);
if (CollectionUtils.isNotEmpty(sameZoneInvokers)) {
// 在同区域内使用Failover策略
return failoverInvoker.doInvoke(invocation, sameZoneInvokers, loadbalance);
}
// 同区域无可用Invoker时,使用所有Invoker
return failoverInvoker.invoke(invocation);
}
}
5. 高级特性
5.1 服务降级
Dubbo支持服务降级功能,当服务调用失败时可以返回默认值或执行降级逻辑:
1
2
3
// 配置方式
<dubbo:reference interface="com.foo.BarService" mock="return null" />
<dubbo:reference interface="com.foo.BarService" mock="com.foo.BarServiceMock" />
5.2 集群扩展点
Dubbo的集群模块提供了多个扩展点,允许用户自定义集群实现:
- Cluster: 自定义集群策略
- LoadBalance: 自定义负载均衡算法
- Router: 自定义路由规则
- Directory: 自定义服务目录
5.3 Dubbo3的服务网格支持
Dubbo3增强了与服务网格的集成能力:
- 应用级服务发现: 支持按应用粒度进行服务注册发现
- Triple协议: 基于HTTP/2的RPC协议,兼容gRPC
- Proxyless模式: 支持无代理服务网格模式
6. 最佳实践
6.1 选择合适的集群策略
- 对于读操作,推荐使用
failover策略 - 对于写操作,推荐使用
failfast策略 - 对于通知操作,推荐使用
failback策略 - 对于多活部署,推荐使用
zone-aware策略
6.2 调优重试参数
1
2
3
4
5
<!-- 配置重试次数 -->
<dubbo:reference interface="com.foo.BarService" retries="2" />
<!-- 配置重试间隔 -->
<dubbo:reference interface="com.foo.BarService" retries-interval="100" />
6.3 配置超时与并发控制
1
2
3
4
5
<!-- 配置调用超时 -->
<dubbo:reference interface="com.foo.BarService" timeout="5000" />
<!-- 配置并发控制 -->
<dubbo:service interface="com.foo.BarService" executes="10" />
7. 总结
Dubbo的集群模块通过灵活的设计,提供了丰富的容错策略,使服务消费者能够在复杂的分布式环境中可靠地调用服务。在Dubbo3中,集群模块进一步增强了对云原生环境的支持,提供了更多的容错策略和服务治理能力。
通过本文的分析,我们深入了解了Dubbo集群模块的设计思想和实现细节,包括Cluster接口、Directory服务目录、各种容错策略的实现原理等。这些知识对于理解Dubbo的服务调用流程和进行服务治理都非常重要。
希望本文能帮助读者更好地理解和使用Dubbo的集群功能,构建更加可靠的分布式系统。

