前面五篇我们搞懂了Dubbo的骨架、血型、开门、找门和通信协议。今天来聊一个最"救命"的功能——集群容错。如果你的服务部署了3台机器,其中一台突然挂了,Dubbo是怎么做到"自动切换"、让用户完全无感知的?这就像一个球队有5个前锋,一个受伤了,教练马上换另一个上,观众根本不知道发生了什么。
一、为什么需要集群容错?
1.1 分布式系统的" Murphy定律"
分布式系统里,任何可能出错的地方,一定会出错:
网络抖动:某个Provider突然连不上了,过5秒又恢复了 机器宕机:一台服务器挂了,服务实例少了一个 服务过载:某台Provider压力太大,响应变慢甚至超时 版本不兼容:升级后某个Provider有问题,需要回滚
如果没有集群容错,Consumer会傻乎乎地一直给挂掉的Provider发请求,直到自己也被拖死。
1.2 Dubbo的"保险策略"
Dubbo提供了6种集群容错策略:
默认是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<T> extendsAbstractClusterInvoker<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的逻辑就像"打地鼠":
先拿到锤子(getMethodParameter获取retries参数) 看地鼠从哪个洞出来(LoadBalance选择Provider) 打一锤子(调用invoker.invoke()) 没打中?换个洞继续(catch异常,重新选择,重试) 打了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<T> extendsAbstractClusterInvoker<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<T> extendsAbstractClusterInvoker<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(null, null, invocation); } }}Failsafe就像"打碎了花瓶,默默扫掉"——适用于日志记录、监控上报这种"失败了也无所谓"的场景。
4.3 FailbackCluster:失败自动恢复
// FailbackClusterInvoker.javapublicclassFailbackClusterInvoker<T> extendsAbstractClusterInvoker<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(null, null, 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<T> extendsAbstractClusterInvoker<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)读取方法级配置。
六、总结:集群容错的"球队战术"
Dubbo的集群容错就像足球队的战术板——Failover是"主力+替补",Failfast是"一击必杀",Failsafe是"输了也不影响大局"。根据业务场景选对策略,才能让分布式系统既高可用、又不出错。
下一篇预告:负载均衡与路由——请求该去哪儿——看Dubbo是怎么在多个Provider之间"公平分配"请求的,以及Router是怎么做灰度发布、蓝绿部署的。
Dubbo 3.3.x 源码深度解读系列(六) | 持续更新中
夜雨聆风