文章

5. Dubbo3 集群源码分析

5. Dubbo3 集群源码分析

集群

参考链接:dubbo2.6源码分析

1. 集群模块概述

在分布式系统中,为了保证高可用性,同一服务通常会部署多个实例。对于服务消费者而言,如何从这些实例中选择一个进行调用,以及如何处理调用失败的情况,是需要特别关注的问题。Dubbo的集群模块正是为解决这些问题而设计的。

集群模块的核心职责包括:

  1. 将多个服务提供者合并为一个虚拟的服务提供点
  2. 在调用失败时,根据配置的策略进行重试或其他容错处理
  3. 屏蔽服务提供者集群的复杂性,简化服务消费者的调用逻辑

Dubbo通过Cluster接口和Cluster Invoker的设计模式实现了这些功能。Cluster负责将多个服务提供者合并为一个Cluster Invoker,而Cluster Invoker则负责具体的调用逻辑和容错策略。

2. 集群容错架构

在分析具体代码之前,我们先了解Dubbo集群容错的整体架构。

集群容错架构

2.1 核心组件

Dubbo集群容错涉及以下核心组件:

  1. Cluster: 集群接口,将多个Invoker合并成一个Invoker
  2. Directory: 服务目录,存储Invoker列表,可动态感知服务提供者变化
  3. Router: 路由器,根据规则筛选Invoker
  4. LoadBalance: 负载均衡器,从多个Invoker中选择一个
  5. Cluster Invoker: 集群Invoker,封装了集群容错逻辑

集群实现

2.2 集群工作流程

集群工作流程分为两个阶段:

  1. 初始化阶段:服务消费者启动时,Cluster实现类为服务消费者创建Cluster Invoker实例。
  2. 调用阶段:服务消费者发起远程调用时,Cluster Invoker执行以下步骤:
    • 调用Directory的list方法获取Invoker列表
    • 调用Router的route方法进行路由,过滤不符合条件的Invoker
    • 使用LoadBalance从Invoker列表中选择一个Invoker
    • 调用选中Invoker的invoke方法进行远程调用
    • 根据配置的容错策略处理调用结果或异常

3. 容错策略详解

Dubbo提供了多种容错策略,以适应不同的业务场景:

容错策略配置值功能描述适用场景
Failover Clusterfailover失败自动切换,重试其他服务器适用于幂等操作,如读取操作
Failfast Clusterfailfast快速失败,只发起一次调用适用于非幂等操作,如新增记录
Failsafe Clusterfailsafe失败安全,出现异常时忽略适用于日志记录等操作
Failback Clusterfailback失败自动恢复,定时重发适用于消息通知等操作
Forking Clusterforking并行调用多个服务提供者适用于实时性要求较高的读操作
Broadcast Clusterbroadcast广播调用所有提供者适用于通知所有提供者更新缓存等操作

Dubbo3新增了以下容错策略:

容错策略配置值功能描述适用场景
ZoneAware Clusterzone-aware区域感知,优先调用同区域服务适用于多区域部署场景
Available Clusteravailable调用第一个可用的服务提供者适用于服务状态查询等非关键操作

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的集群模块提供了多个扩展点,允许用户自定义集群实现:

  1. Cluster: 自定义集群策略
  2. LoadBalance: 自定义负载均衡算法
  3. Router: 自定义路由规则
  4. Directory: 自定义服务目录

5.3 Dubbo3的服务网格支持

Dubbo3增强了与服务网格的集成能力:

  1. 应用级服务发现: 支持按应用粒度进行服务注册发现
  2. Triple协议: 基于HTTP/2的RPC协议,兼容gRPC
  3. 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的集群功能,构建更加可靠的分布式系统。

本文由作者按照 CC BY 4.0 进行授权