乐于分享
好东西不私藏

Dubbo 3.3.x 源码深度解读(六):集群容错——从Failover到ClusterInvoker

Dubbo 3.3.x 源码深度解读(六):集群容错——从Failover到ClusterInvoker

前面五篇我们搞懂了Dubbo的骨架、血型、开门、找门和通信协议。今天来聊一个最"救命"的功能——集群容错。如果你的服务部署了3台机器,其中一台突然挂了,Dubbo是怎么做到"自动切换"、让用户完全无感知的?这就像一个球队有5个前锋,一个受伤了,教练马上换另一个上,观众根本不知道发生了什么。

一、为什么需要集群容错?

1.1 分布式系统的" Murphy定律"

分布式系统里,任何可能出错的地方,一定会出错:

  • 网络抖动:某个Provider突然连不上了,过5秒又恢复了
  • 机器宕机:一台服务器挂了,服务实例少了一个
  • 服务过载:某台Provider压力太大,响应变慢甚至超时
  • 版本不兼容:升级后某个Provider有问题,需要回滚

如果没有集群容错,Consumer会傻乎乎地一直给挂掉的Provider发请求,直到自己也被拖死。

1.2 Dubbo的"保险策略"

Dubbo提供了6种集群容错策略:

策略
名称
行为
适用场景
Failover
失败自动切换
失败换下一个Provider重试
读操作、幂等操作
Failfast
快速失败
只调用一次,失败立即抛异常
非幂等写操作
Failsafe
失败安全
失败忽略,返回空结果
日志记录、监控上报
Failback
失败自动恢复
失败异步重试,定时任务补发
消息通知类
Forking
并行调用
同时调用多个Provider,取最快响应
实时性要求高的场景
Broadcast
广播调用
逐个调用所有Provider
广播通知、缓存刷新

默认是Failover。今天主要讲这个"顶梁柱"。


二、Cluster接口:容错策略的入口

2.1 Cluster的SPI定义

// org.apache.dubbo.rpc.cluster.Cluster@SPI(FailoverCluster.NAME)  // 默认FailoverpublicinterfaceCluster{@Adaptive    <T> Invoker<T> join(Directory<T> directory)throws RpcException;}

Cluster是一个SPI扩展点,通过cluster=failover(或failfast、failsafe等)参数切换实现。

2.2 Cluster的6种实现

// 6种Cluster实现publicclassFailoverClusterimplementsCluster{}    // 失败重试publicclassFailfastClusterimplementsCluster{}     // 快速失败publicclassFailsafeClusterimplementsCluster{}     // 失败安全publicclassFailbackClusterimplementsCluster{}     // 失败恢复publicclassForkingClusterimplementsCluster{}       // 并行调用publicclassBroadcastClusterimplementsCluster{}     // 广播调用

每种Cluster实现,只干一件事:根据Directory里的Provider列表,返回一个包装后的Invoker。真正的容错逻辑,在返回的ClusterInvoker里。


三、FailoverCluster:失败自动切换

3.1 FailoverCluster.join()

// FailoverCluster.javapublicclassFailoverClusterimplementsCluster{publicfinalstatic String NAME = "failover";@Overridepublic <T> Invoker<T> join(Directory<T> directory)throws RpcException {// 返回FailoverClusterInvoker,包装了Directoryreturnnew FailoverClusterInvoker<>(directory);    }}

FailoverCluster就像"替补席教练"——它不直接上场踢球,而是安排FailoverClusterInvoker上场,并告诉它"失败了换下一个"。

3.2 FailoverClusterInvoker:真正的容错逻辑

// FailoverClusterInvoker.javapublicclassFailoverClusterInvoker<TextendsAbstractClusterInvoker<T{@Overridepublic Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)throws RpcException {// 第1步:获取重试次数// retries参数默认是2,所以len = 2 + 1 = 3(总共调用3次)int len = getUrl().getMethodParameter(            invocation.getMethodName(),             RETRIES_KEY,             DEFAULT_RETRIES  // 默认值2        ) + 1;// 记录最后一次异常        RpcException le = null;// 记录已经调用过的Provider,避免重复调用同一个        List<Invoker<T>> invoked = new ArrayList<Invoker<T>>(invokers.size());// 第2步:循环调用,失败就重试for (int i = 0; i < len; i++) {// 如果已经调用过一些Provider了,重新获取Provider列表(可能变了)if (i > 0) {                checkWheatherDestoried();  // 检查自己是否被销毁了// 重新获取Provider列表// 这里会触发Directory.list(),获取最新的Provider                copyinvokers = list(invocation);// 重新检查                checkInvokers(copyinvokers, invocation);            }// 第3步:用LoadBalance选择一个Provider// select方法内部会做"排除已失败"的逻辑            Invoker<T> invoker = select(loadbalance, invocation, copyinvokers, invoked);// 记录这个Provider已经被调用过了            invoked.add(invoker);// 第4步:真正发起远程调用try {                Result result = invoker.invoke(invocation);return result;  // 成功了!返回结果            } catch (RpcException e) {if (e.isBiz()) {// 业务异常(不是系统异常),直接抛出去throw e;                }// 记录异常,继续重试                le = e;// 如果异常是不可重试的(如参数错误),直接抛出去if (!e.isRetriable()) {throw e;                }            } catch (Throwable e) {                le = new RpcException(e.getMessage(), e);            }        }// 第5步:所有重试都失败了thrownew RpcException("Failed to invoke the method " + invocation.getMethodName()            + " in the service " + getInterface().getName()            + ". Tried " + len + " times...",            le        );    }}

FailoverClusterInvoker的逻辑就像"打地鼠":

  1. 先拿到锤子(getMethodParameter获取retries参数)
  2. 看地鼠从哪个洞出来(LoadBalance选择Provider)
  3. 打一锤子(调用invoker.invoke())
  4. 没打中?换个洞继续(catch异常,重新选择,重试)
  5. 打了3次都没中?不玩了(抛出RpcException)

3.3 AbstractClusterInvoker.select():排除已失败的Provider

// AbstractClusterInvoker.javaprotected Invoker<T> select(LoadBalance loadbalance, Invocation invocation, List<Invoker<T>> invokers, List<Invoker<T>> selected)throws RpcException {if (invokers == null || invokers.isEmpty()) {returnnull;    }    String methodName = invocation.getMethodName();// 获取sticky配置(粘滞连接)boolean sticky = getUrl().getMethodParameter(methodName, CLUSTER_STICKY_KEY, DEFAULT_CLUSTER_STICKY);// 如果启用了sticky,尽量调用同一个Provider(除非它挂了)if (sticky && selected != null && !selected.isEmpty()) {// ... 粘滞逻辑    }// 第1步:用LoadBalance选择一个Provider    Invoker<T> invoker = doSelect(loadbalance, invocation, invokers, selected);// 第2步:如果选中的Provider在"已失败"列表里,排除它重新选if (selected != null && selected.contains(invoker)) {// 重新选择,排除已失败的        invoker = reselect(loadbalance, invocation, invokers, selected);    }return invoker;}

AbstractClusterInvoker.select方法的Provider选择逻辑。select方法在LoadBalance选择的基础上增加了两个优化:粘滞连接(sticky)和已失败Provider排除。sticky配置使Consumer尽量调用同一个Provider(利用连接复用和缓存),除非该Provider已失败。doSelect调用LoadBalance(如RandomLoadBalance)选择一个Provider,如果选中的Provider在selected列表(本次调用已失败过的Provider列表)中,则调用reselect重新选择排除已失败的。

这种设计在Failover场景下特别重要——第一次调用Provider A失败后重试时,select会排除A选择其他Provider,避免再次调用已失败的节点。selected列表在ClusterInvoker的循环重试中维护,每次重试将失败的Provider加入selected,确保后续重试不会选到同一个失败节点。

select()方法的关键:

  • LoadBalance选:先让负载均衡算法选一个
  • 排除已失败:如果这个已经被调用过了(在selected列表里),重新选
  • 粘滞连接:如果配置了sticky,尽量一直调用同一个Provider(适用于有状态服务)

四、其他集群容错策略速览

4.1 FailfastCluster:快速失败

// FailfastClusterInvoker.javapublicclassFailfastClusterInvoker<TextendsAbstractClusterInvoker<T{@Overridepublic Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)throws RpcException {// 只调用一次,失败立即抛异常        checkInvokers(invokers, invocation);        Invoker<T> invoker = select(loadbalance, invocation, invokers, null);try {return invoker.invoke(invocation);        } catch (Throwable e) {// 失败就抛异常,不 retryif (e instanceof RpcException) {throw (RpcException) e;            }thrownew RpcException(e);        }    }}

Failfast就像"一击不中,立即报警"——适用于非幂等的写操作(如扣款、下单),失败了不能重试。

4.2 FailsafeCluster:失败安全

// FailsafeClusterInvoker.javapublicclassFailsafeClusterInvoker<TextendsAbstractClusterInvoker<T{@Overridepublic Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)throws RpcException {try {            checkInvokers(invokers, invocation);            Invoker<T> invoker = select(loadbalance, invocation, invokers, null);return invoker.invoke(invocation);        } catch (Throwable e) {// 失败了?忽略,打印个日志,返回空结果            logger.error("Failsafe ignore exception: " + e.getMessage(), e);return AsyncRpcResult.newDefaultAsyncResult(nullnull, invocation);        }    }}

Failsafe就像"打碎了花瓶,默默扫掉"——适用于日志记录、监控上报这种"失败了也无所谓"的场景。

4.3 FailbackCluster:失败自动恢复

// FailbackClusterInvoker.javapublicclassFailbackClusterInvoker<TextendsAbstractClusterInvoker<T{// 失败的任务队列privatefinal ConcurrentMap<Invocation, AbstractClusterInvoker<?>> failed = new ConcurrentHashMap<>();@Overridepublic Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)throws RpcException {try {            checkInvokers(invokers, invocation);            Invoker<T> invoker = select(loadbalance, invocation, invokers, null);return invoker.invoke(invocation);        } catch (Throwable e) {// 失败了?放到队列里,稍后异步重试            addFailed(invocation, invokers, loadbalance);return AsyncRpcResult.newDefaultAsyncResult(nullnull, invocation);        }    }// 定时任务:每隔5秒重试失败的任务privatevoidretryFailed(){for (Map.Entry<Invocation, AbstractClusterInvoker<?>> entry : failed.entrySet()) {try {                entry.getValue().invoke(entry.getKey());                failed.remove(entry.getKey());  // 重试成功,移除            } catch (Throwable e) {// 重试又失败了,继续留在队列里,下次再试            }        }    }}

Failback就像"快递送不到?先放驿站,明天再送"——适用于消息通知类场景,失败了不阻塞主流程,异步补偿重试。

4.4 ForkingCluster:并行调用

// ForkingClusterInvoker.javapublicclassForkingClusterInvoker<TextendsAbstractClusterInvoker<T{@Overridepublic Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance)throws RpcException {// 获取并行调用数量(默认2个)int forks = getUrl().getParameter(FORKS_KEY, DEFAULT_FORKS);// 用线程池并行调用多个Provider        ExecutorService executor = getExecutor();        CountDownLatch countDownLatch = new CountDownLatch(forks);        AtomicReference<Object> resultRef = new AtomicReference<>();for (Invoker<T> invoker : selected) {            executor.execute(() -> {try {                    Result result = invoker.invoke(invocation);                    resultRef.set(result.getValue());                    countDownLatch.countDown();                } catch (Throwable e) {                    countDownLatch.countDown();                }            });        }// 等第一个返回的结果        countDownLatch.await(timeout, TimeUnit.MILLISECONDS);return (Result) resultRef.get();    }}

Forking就像"同时叫3辆出租车,谁先到坐谁的"——适用于实时性要求高的场景(如金融行情查询)。


五、集群容错的配置与使用

5.1 配置方式

<!-- XML配置 --><dubbo:referenceid="demoService"interface="com.example.DemoService"cluster="failover"retries="3" /><dubbo:referenceid="orderService"interface="com.example.OrderService"cluster="failfast" /><dubbo:referenceid="logService"interface="com.example.LogService"cluster="failsafe" />
// 注解配置@DubboReference(cluster = "failover", retries = 3)private DemoService demoService;@DubboReference(cluster = "failfast")private OrderService orderService;

集群容错策略的注解配置方式。@DubboReference注解的cluster属性指定容错策略——failover(失败自动切换,默认策略,自动重试切换到其他Provider)适合幂等性操作如查询;failfast(快速失败,不重试直接抛异常)适合非幂等性操作如新增记录,防止重复执行导致数据不一致。retries属性配置重试次数,failover策略默认重试2次(共3次调用)。这些配置通过URL传递到Cluster层,ClusterInvoker根据cluster参数通过SPI加载对应的Cluster实现(FailoverCluster、FailfastCluster等)。Dubbo支持的容错策略还包括failsafe(失败忽略,记录日志不抛异常)、failback(失败后台重试)、forking(并行调用多个Provider)、broadcast(广播调用所有Provider)。

// API配置ReferenceConfig<DemoService> reference = new ReferenceConfig<>();reference.setInterface(DemoService.class);reference.setCluster("failover");reference.setRetries(3);

集群容错的API配置方式。ReferenceConfig是Dubbo的编程式配置API,等效于@DubboReference注解。setCluster设置容错策略为failover,setRetries设置重试次数为3。API配置适用于动态创建Reference的场景,如网关服务需要根据请求参数动态引用不同的Dubbo服务。与注解配置相比,API配置更灵活但更繁琐,需要手动管理ReferenceConfig的生命周期(初始化和销毁)。ReferenceConfig内部实现了懒加载——第一次调用get()方法时才真正创建Invoker并连接注册中心。配置完成后通过get()获取服务代理对象,后续调用代理方法时触发远程调用。在生产环境中,API配置常与Spring的@Conditional结合实现服务的动态引用。

5.2 配置优先级

Consumer端方法级 > Consumer端接口级 > Provider端方法级 > Provider端接口级 > Consumer端全局 > Provider端全局 > 默认值

比如你可以在Consumer端为某个方法单独配置:

@DubboReference(    cluster = "failover",    retries = 3,    parameters = {"sayHello.retries""5"}  // sayHello方法单独配置5次重试)private DemoService demoService;

Dubbo的配置优先级机制。@DubboReference的全局配置cluster=failover和retries=3应用于DemoService的所有方法。parameters属性支持方法级别的细粒度配置——{"sayHello.retries", "5"}表示sayHello方法单独配置5次重试,覆盖全局的3次。Dubbo的配置优先级从高到低为:方法级配置 > 接口级配置 > 全局配置 > 默认值。这种层级配置使得不同方法可以有差异化的容错策略——查询方法可以多重试(retries=5),写入方法少重试或快速失败(retries=0配合failfast)。配置在URL中以参数形式传递,如sayHello.retries=5编码为URL参数,ClusterInvoker在调用时通过URL.getMethodParameter(methodName, "retries", default)读取方法级配置。


六、总结:集群容错的"球队战术"

策略
战术比喻
适用场景
源码实现
Failover
前锋受伤了换替补
读操作、幂等操作
FailoverClusterInvoker,循环重试
Failfast
一击不中立即认输
非幂等写操作(扣款、下单)
FailfastClusterInvoker,直接抛异常
Failsafe
打碎了花瓶扫掉假装没看见
日志、监控上报
FailsafeClusterInvoker,忽略异常返回空
Failback
快递送不到明天再送
消息通知
FailbackClusterInvoker,异步队列重试
Forking
同时叫3辆出租车
实时性要求高
ForkingClusterInvoker,线程池并行调用
Broadcast
逐个通知全班同学
广播刷新缓存
BroadcastClusterInvoker,逐个调用

Dubbo的集群容错就像足球队的战术板——Failover是"主力+替补",Failfast是"一击必杀",Failsafe是"输了也不影响大局"。根据业务场景选对策略,才能让分布式系统既高可用、又不出错。

下一篇预告:负载均衡与路由——请求该去哪儿——看Dubbo是怎么在多个Provider之间"公平分配"请求的,以及Router是怎么做灰度发布、蓝绿部署的。


Dubbo 3.3.x 源码深度解读系列(六) | 持续更新中