3.2 Cluster Invoker 分析

我们首先从各种 Cluster Invoker 的父类 AbstractClusterInvoker 源码开始说起。前面说过,集群工作过程可分为两个阶段,第一个阶段是在服务消费者初始化期间,这个在服务引用那篇文章中分析过,就不赘述。第二个阶段是在服务消费者进行远程调用时,此时 AbstractClusterInvoker 的 invoke 方法会被调用。列举 Invoker,负载均衡等操作均会在此阶段被执行。因此下面先来看一下 invoke 方法的逻辑。

  1. public Result invoke(final Invocation invocation) throws RpcException {
  2. checkWhetherDestroyed();
  3. LoadBalance loadbalance = null;
  4. // 绑定 attachments 到 invocation 中.
  5. Map<String, String> contextAttachments = RpcContext.getContext().getAttachments();
  6. if (contextAttachments != null && contextAttachments.size() != 0) {
  7. ((RpcInvocation) invocation).addAttachments(contextAttachments);
  8. }
  9. // 列举 Invoker
  10. List<Invoker<T>> invokers = list(invocation);
  11. if (invokers != null && !invokers.isEmpty()) {
  12. // 加载 LoadBalance
  13. loadbalance = ExtensionLoader.getExtensionLoader(LoadBalance.class).getExtension(invokers.get(0).getUrl()
  14. .getMethodParameter(RpcUtils.getMethodName(invocation), Constants.LOADBALANCE_KEY, Constants.DEFAULT_LOADBALANCE));
  15. }
  16. RpcUtils.attachInvocationIdIfAsync(getUrl(), invocation);
  17. // 调用 doInvoke 进行后续操作
  18. return doInvoke(invocation, invokers, loadbalance);
  19. }
  20. // 抽象方法,由子类实现
  21. protected abstract Result doInvoke(Invocation invocation, List<Invoker<T>> invokers,
  22. LoadBalance loadbalance) throws RpcException;

AbstractClusterInvoker 的 invoke 方法主要用于列举 Invoker,以及加载 LoadBalance。最后再调用模板方法 doInvoke 进行后续操作。下面我们来看一下 Invoker 列举方法 list(Invocation) 的逻辑,如下:

  1. protected List<Invoker<T>> list(Invocation invocation) throws RpcException {
  2. // 调用 Directory 的 list 方法列举 Invoker
  3. List<Invoker<T>> invokers = directory.list(invocation);
  4. return invokers;
  5. }

如上,AbstractClusterInvoker 中的 list 方法做的事情很简单,只是简单的调用了 Directory 的 list 方法,没有其他更多的逻辑了。Directory 即相关实现类在前文已经分析过,这里就不多说了。接下来,我们把目光转移到 AbstractClusterInvoker 的各种实现类上,来看一下这些实现类是如何实现 doInvoke 方法逻辑的。

3.2.1 FailoverClusterInvoker

FailoverClusterInvoker 在调用失败时,会自动切换 Invoker 进行重试。默认配置下,Dubbo 会使用这个类作为缺省 Cluster Invoker。下面来看一下该类的逻辑。

  1. public class FailoverClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. // 省略部分代码
  3. @Override
  4. public Result doInvoke(Invocation invocation, final List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  5. List<Invoker<T>> copyinvokers = invokers;
  6. checkInvokers(copyinvokers, invocation);
  7. // 获取重试次数
  8. int len = getUrl().getMethodParameter(invocation.getMethodName(), Constants.RETRIES_KEY, Constants.DEFAULT_RETRIES) + 1;
  9. if (len <= 0) {
  10. len = 1;
  11. }
  12. RpcException le = null;
  13. List<Invoker<T>> invoked = new ArrayList<Invoker<T>>(copyinvokers.size());
  14. Set<String> providers = new HashSet<String>(len);
  15. // 循环调用,失败重试
  16. for (int i = 0; i < len; i++) {
  17. if (i > 0) {
  18. checkWhetherDestroyed();
  19. // 在进行重试前重新列举 Invoker,这样做的好处是,如果某个服务挂了,
  20. // 通过调用 list 可得到最新可用的 Invoker 列表
  21. copyinvokers = list(invocation);
  22. // 对 copyinvokers 进行判空检查
  23. checkInvokers(copyinvokers, invocation);
  24. }
  25. // 通过负载均衡选择 Invoker
  26. Invoker<T> invoker = select(loadbalance, invocation, copyinvokers, invoked);
  27. // 添加到 invoker 到 invoked 列表中
  28. invoked.add(invoker);
  29. // 设置 invoked 到 RPC 上下文中
  30. RpcContext.getContext().setInvokers((List) invoked);
  31. try {
  32. // 调用目标 Invoker 的 invoke 方法
  33. Result result = invoker.invoke(invocation);
  34. return result;
  35. } catch (RpcException e) {
  36. if (e.isBiz()) {
  37. throw e;
  38. }
  39. le = e;
  40. } catch (Throwable e) {
  41. le = new RpcException(e.getMessage(), e);
  42. } finally {
  43. providers.add(invoker.getUrl().getAddress());
  44. }
  45. }
  46. // 若重试失败,则抛出异常
  47. throw new RpcException(..., "Failed to invoke the method ...");
  48. }
  49. }

如上,FailoverClusterInvoker 的 doInvoke 方法首先是获取重试次数,然后根据重试次数进行循环调用,失败后进行重试。在 for 循环内,首先是通过负载均衡组件选择一个 Invoker,然后再通过这个 Invoker 的 invoke 方法进行远程调用。如果失败了,记录下异常,并进行重试。重试时会再次调用父类的 list 方法列举 Invoker。整个流程大致如此,不是很难理解。下面我们看一下 select 方法的逻辑。

  1. protected Invoker<T> select(LoadBalance loadbalance, Invocation invocation, List<Invoker<T>> invokers, List<Invoker<T>> selected) throws RpcException {
  2. if (invokers == null || invokers.isEmpty())
  3. return null;
  4. // 获取调用方法名
  5. String methodName = invocation == null ? "" : invocation.getMethodName();
  6. // 获取 sticky 配置,sticky 表示粘滞连接。所谓粘滞连接是指让服务消费者尽可能的
  7. // 调用同一个服务提供者,除非该提供者挂了再进行切换
  8. boolean sticky = invokers.get(0).getUrl().getMethodParameter(methodName, Constants.CLUSTER_STICKY_KEY, Constants.DEFAULT_CLUSTER_STICKY);
  9. {
  10. // 检测 invokers 列表是否包含 stickyInvoker,如果不包含,
  11. // 说明 stickyInvoker 代表的服务提供者挂了,此时需要将其置空
  12. if (stickyInvoker != null && !invokers.contains(stickyInvoker)) {
  13. stickyInvoker = null;
  14. }
  15. // 在 sticky 为 true,且 stickyInvoker != null 的情况下。如果 selected 包含
  16. // stickyInvoker,表明 stickyInvoker 对应的服务提供者可能因网络原因未能成功提供服务。
  17. // 但是该提供者并没挂,此时 invokers 列表中仍存在该服务提供者对应的 Invoker。
  18. if (sticky && stickyInvoker != null && (selected == null || !selected.contains(stickyInvoker))) {
  19. // availablecheck 表示是否开启了可用性检查,如果开启了,则调用 stickyInvoker 的
  20. // isAvailable 方法进行检查,如果检查通过,则直接返回 stickyInvoker。
  21. if (availablecheck && stickyInvoker.isAvailable()) {
  22. return stickyInvoker;
  23. }
  24. }
  25. }
  26. // 如果线程走到当前代码处,说明前面的 stickyInvoker 为空,或者不可用。
  27. // 此时继续调用 doSelect 选择 Invoker
  28. Invoker<T> invoker = doSelect(loadbalance, invocation, invokers, selected);
  29. // 如果 sticky 为 true,则将负载均衡组件选出的 Invoker 赋值给 stickyInvoker
  30. if (sticky) {
  31. stickyInvoker = invoker;
  32. }
  33. return invoker;
  34. }

如上,select 方法的主要逻辑集中在了对粘滞连接特性的支持上。首先是获取 sticky 配置,然后再检测 invokers 列表中是否包含 stickyInvoker,如果不包含,则认为该 stickyInvoker 不可用,此时将其置空。这里的 invokers 列表可以看做是存活着的服务提供者列表,如果这个列表不包含 stickyInvoker,那自然而然的认为 stickyInvoker 挂了,所以置空。如果 stickyInvoker 存在于 invokers 列表中,此时要进行下一项检测 — 检测 selected 中是否包含 stickyInvoker。如果包含的话,说明 stickyInvoker 在此之前没有成功提供服务(但其仍然处于存活状态)。此时我们认为这个服务不可靠,不应该在重试期间内再次被调用,因此这个时候不会返回该 stickyInvoker。如果 selected 不包含 stickyInvoker,此时还需要进行可用性检测,比如检测服务提供者网络连通性等。当可用性检测通过,才可返回 stickyInvoker,否则调用 doSelect 方法选择 Invoker。如果 sticky 为 true,此时会将 doSelect 方法选出的 Invoker 赋值给 stickyInvoker。

以上就是 select 方法的逻辑,这段逻辑看起来不是很复杂,但是信息量比较大。不搞懂 invokers 和 selected 两个入参的含义,以及粘滞连接特性,这段代码是不容易看懂的。所以大家在阅读这段代码时,不要忽略了对背景知识的理解。关于 select 方法先分析这么多,继续向下分析。

  1. private Invoker<T> doSelect(LoadBalance loadbalance, Invocation invocation, List<Invoker<T>> invokers, List<Invoker<T>> selected) throws RpcException {
  2. if (invokers == null || invokers.isEmpty())
  3. return null;
  4. if (invokers.size() == 1)
  5. return invokers.get(0);
  6. if (loadbalance == null) {
  7. // 如果 loadbalance 为空,这里通过 SPI 加载 Loadbalance,默认为 RandomLoadBalance
  8. loadbalance = ExtensionLoader.getExtensionLoader(LoadBalance.class).getExtension(Constants.DEFAULT_LOADBALANCE);
  9. }
  10. // 通过负载均衡组件选择 Invoker
  11. Invoker<T> invoker = loadbalance.select(invokers, getUrl(), invocation);
  12. // 如果 selected 包含负载均衡选择出的 Invoker,或者该 Invoker 无法经过可用性检查,此时进行重选
  13. if ((selected != null && selected.contains(invoker))
  14. || (!invoker.isAvailable() && getUrl() != null && availablecheck)) {
  15. try {
  16. // 进行重选
  17. Invoker<T> rinvoker = reselect(loadbalance, invocation, invokers, selected, availablecheck);
  18. if (rinvoker != null) {
  19. // 如果 rinvoker 不为空,则将其赋值给 invoker
  20. invoker = rinvoker;
  21. } else {
  22. // rinvoker 为空,定位 invoker 在 invokers 中的位置
  23. int index = invokers.indexOf(invoker);
  24. try {
  25. // 获取 index + 1 位置处的 Invoker,以下代码等价于:
  26. // invoker = invokers.get((index + 1) % invokers.size());
  27. invoker = index < invokers.size() - 1 ? invokers.get(index + 1) : invokers.get(0);
  28. } catch (Exception e) {
  29. logger.warn("... may because invokers list dynamic change, ignore.");
  30. }
  31. }
  32. } catch (Throwable t) {
  33. logger.error("cluster reselect fail reason is : ...");
  34. }
  35. }
  36. return invoker;
  37. }

doSelect 主要做了两件事,第一是通过负载均衡组件选择 Invoker。第二是,如果选出来的 Invoker 不稳定,或不可用,此时需要调用 reselect 方法进行重选。若 reselect 选出来的 Invoker 为空,此时定位 invoker 在 invokers 列表中的位置 index,然后获取 index + 1 处的 invoker,这也可以看做是重选逻辑的一部分。下面我们来看一下 reselect 方法的逻辑。

  1. private Invoker<T> reselect(LoadBalance loadbalance, Invocation invocation,
  2. List<Invoker<T>> invokers, List<Invoker<T>> selected, boolean availablecheck) throws RpcException {
  3. List<Invoker<T>> reselectInvokers = new ArrayList<Invoker<T>>(invokers.size() > 1 ? (invokers.size() - 1) : invokers.size());
  4. // 下面的 if-else 分支逻辑有些冗余,pull request #2826 对这段代码进行了简化,可以参考一下
  5. // 根据 availablecheck 进行不同的处理
  6. if (availablecheck) {
  7. // 遍历 invokers 列表
  8. for (Invoker<T> invoker : invokers) {
  9. // 检测可用性
  10. if (invoker.isAvailable()) {
  11. // 如果 selected 列表不包含当前 invoker,则将其添加到 reselectInvokers 中
  12. if (selected == null || !selected.contains(invoker)) {
  13. reselectInvokers.add(invoker);
  14. }
  15. }
  16. }
  17. // reselectInvokers 不为空,此时通过负载均衡组件进行选择
  18. if (!reselectInvokers.isEmpty()) {
  19. return loadbalance.select(reselectInvokers, getUrl(), invocation);
  20. }
  21. // 不检查 Invoker 可用性
  22. } else {
  23. for (Invoker<T> invoker : invokers) {
  24. // 如果 selected 列表不包含当前 invoker,则将其添加到 reselectInvokers 中
  25. if (selected == null || !selected.contains(invoker)) {
  26. reselectInvokers.add(invoker);
  27. }
  28. }
  29. if (!reselectInvokers.isEmpty()) {
  30. // 通过负载均衡组件进行选择
  31. return loadbalance.select(reselectInvokers, getUrl(), invocation);
  32. }
  33. }
  34. {
  35. // 若线程走到此处,说明 reselectInvokers 集合为空,此时不会调用负载均衡组件进行筛选。
  36. // 这里从 selected 列表中查找可用的 Invoker,并将其添加到 reselectInvokers 集合中
  37. if (selected != null) {
  38. for (Invoker<T> invoker : selected) {
  39. if ((invoker.isAvailable())
  40. && !reselectInvokers.contains(invoker)) {
  41. reselectInvokers.add(invoker);
  42. }
  43. }
  44. }
  45. if (!reselectInvokers.isEmpty()) {
  46. // 再次进行选择,并返回选择结果
  47. return loadbalance.select(reselectInvokers, getUrl(), invocation);
  48. }
  49. }
  50. return null;
  51. }

reselect 方法总结下来其实只做了两件事情,第一是查找可用的 Invoker,并将其添加到 reselectInvokers 集合中。第二,如果 reselectInvokers 不为空,则通过负载均衡组件再次进行选择。其中第一件事情又可进行细分,一开始,reselect 从 invokers 列表中查找有效可用的 Invoker,若未能找到,此时再到 selected 列表中继续查找。关于 reselect 方法就先分析到这,继续分析其他的 Cluster Invoker。

3.2.2 FailbackClusterInvoker

FailbackClusterInvoker 会在调用失败后,返回一个空结果给服务消费者。并通过定时任务对失败的调用进行重传,适合执行消息通知等操作。下面来看一下它的实现逻辑。

  1. public class FailbackClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. private static final long RETRY_FAILED_PERIOD = 5 * 1000;
  3. private final ScheduledExecutorService scheduledExecutorService = Executors.newScheduledThreadPool(2,
  4. new NamedInternalThreadFactory("failback-cluster-timer", true));
  5. private final ConcurrentMap<Invocation, AbstractClusterInvoker<?>> failed = new ConcurrentHashMap<Invocation, AbstractClusterInvoker<?>>();
  6. private volatile ScheduledFuture<?> retryFuture;
  7. @Override
  8. protected Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  9. try {
  10. checkInvokers(invokers, invocation);
  11. // 选择 Invoker
  12. Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
  13. // 进行调用
  14. return invoker.invoke(invocation);
  15. } catch (Throwable e) {
  16. // 如果调用过程中发生异常,此时仅打印错误日志,不抛出异常
  17. logger.error("Failback to invoke method ...");
  18. // 记录调用信息
  19. addFailed(invocation, this);
  20. // 返回一个空结果给服务消费者
  21. return new RpcResult();
  22. }
  23. }
  24. private void addFailed(Invocation invocation, AbstractClusterInvoker<?> router) {
  25. if (retryFuture == null) {
  26. synchronized (this) {
  27. if (retryFuture == null) {
  28. // 创建定时任务,每隔5秒执行一次
  29. retryFuture = scheduledExecutorService.scheduleWithFixedDelay(new Runnable() {
  30. @Override
  31. public void run() {
  32. try {
  33. // 对失败的调用进行重试
  34. retryFailed();
  35. } catch (Throwable t) {
  36. // 如果发生异常,仅打印异常日志,不抛出
  37. logger.error("Unexpected error occur at collect statistic", t);
  38. }
  39. }
  40. }, RETRY_FAILED_PERIOD, RETRY_FAILED_PERIOD, TimeUnit.MILLISECONDS);
  41. }
  42. }
  43. }
  44. // 添加 invocation 和 invoker 到 failed 中
  45. failed.put(invocation, router);
  46. }
  47. void retryFailed() {
  48. if (failed.size() == 0) {
  49. return;
  50. }
  51. // 遍历 failed,对失败的调用进行重试
  52. for (Map.Entry<Invocation, AbstractClusterInvoker<?>> entry : new HashMap<Invocation, AbstractClusterInvoker<?>>(failed).entrySet()) {
  53. Invocation invocation = entry.getKey();
  54. Invoker<?> invoker = entry.getValue();
  55. try {
  56. // 再次进行调用
  57. invoker.invoke(invocation);
  58. // 调用成功后,从 failed 中移除 invoker
  59. failed.remove(invocation);
  60. } catch (Throwable e) {
  61. // 仅打印异常,不抛出
  62. logger.error("Failed retry to invoke method ...");
  63. }
  64. }
  65. }
  66. }

这个类主要由3个方法组成,首先是 doInvoker,该方法负责初次的远程调用。若远程调用失败,则通过 addFailed 方法将调用信息存入到 failed 中,等待定时重试。addFailed 在开始阶段会根据 retryFuture 为空与否,来决定是否开启定时任务。retryFailed 方法则是包含了失败重试的逻辑,该方法会对 failed 进行遍历,然后依次对 Invoker 进行调用。调用成功则将 Invoker 从 failed 中移除,调用失败则忽略失败原因。

以上就是 FailbackClusterInvoker 的执行逻辑,不是很复杂,继续往下看。

3.2.3 FailfastClusterInvoker

FailfastClusterInvoker 只会进行一次调用,失败后立即抛出异常。适用于幂等操作,比如新增记录。源码如下:

  1. public class FailfastClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. @Override
  3. public Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  4. checkInvokers(invokers, invocation);
  5. // 选择 Invoker
  6. Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
  7. try {
  8. // 调用 Invoker
  9. return invoker.invoke(invocation);
  10. } catch (Throwable e) {
  11. if (e instanceof RpcException && ((RpcException) e).isBiz()) {
  12. // 抛出异常
  13. throw (RpcException) e;
  14. }
  15. // 抛出异常
  16. throw new RpcException(..., "Failfast invoke providers ...");
  17. }
  18. }
  19. }

如上,首先是通过 select 方法选择 Invoker,然后进行远程调用。如果调用失败,则立即抛出异常。FailfastClusterInvoker 就先分析到这,下面分析 FailsafeClusterInvoker。

3.2.4 FailsafeClusterInvoker

FailsafeClusterInvoker 是一种失败安全的 Cluster Invoker。所谓的失败安全是指,当调用过程中出现异常时,FailsafeClusterInvoker 仅会打印异常,而不会抛出异常。适用于写入审计日志等操作。下面分析源码。

  1. public class FailsafeClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. @Override
  3. public Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  4. try {
  5. checkInvokers(invokers, invocation);
  6. // 选择 Invoker
  7. Invoker<T> invoker = select(loadbalance, invocation, invokers, null);
  8. // 进行远程调用
  9. return invoker.invoke(invocation);
  10. } catch (Throwable e) {
  11. // 打印错误日志,但不抛出
  12. logger.error("Failsafe ignore exception: " + e.getMessage(), e);
  13. // 返回空结果忽略错误
  14. return new RpcResult();
  15. }
  16. }
  17. }

FailsafeClusterInvoker 的逻辑和 FailfastClusterInvoker 的逻辑一样简单,无需过多说明。继续向下分析。

3.2.5 ForkingClusterInvoker

ForkingClusterInvoker 会在运行时通过线程池创建多个线程,并发调用多个服务提供者。只要有一个服务提供者成功返回了结果,doInvoke 方法就会立即结束运行。ForkingClusterInvoker 的应用场景是在一些对实时性要求比较高读操作(注意是读操作,并行写操作可能不安全)下使用,但这将会耗费更多的资源。下面来看该类的实现。

  1. public class ForkingClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. private final ExecutorService executor = Executors.newCachedThreadPool(
  3. new NamedInternalThreadFactory("forking-cluster-timer", true));
  4. @Override
  5. public Result doInvoke(final Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  6. try {
  7. checkInvokers(invokers, invocation);
  8. final List<Invoker<T>> selected;
  9. // 获取 forks 配置
  10. final int forks = getUrl().getParameter(Constants.FORKS_KEY, Constants.DEFAULT_FORKS);
  11. // 获取超时配置
  12. final int timeout = getUrl().getParameter(Constants.TIMEOUT_KEY, Constants.DEFAULT_TIMEOUT);
  13. // 如果 forks 配置不合理,则直接将 invokers 赋值给 selected
  14. if (forks <= 0 || forks >= invokers.size()) {
  15. selected = invokers;
  16. } else {
  17. selected = new ArrayList<Invoker<T>>();
  18. // 循环选出 forks 个 Invoker,并添加到 selected 中
  19. for (int i = 0; i < forks; i++) {
  20. // 选择 Invoker
  21. Invoker<T> invoker = select(loadbalance, invocation, invokers, selected);
  22. if (!selected.contains(invoker)) {
  23. selected.add(invoker);
  24. }
  25. }
  26. }
  27. // ----------------------✨ 分割线1 ✨---------------------- //
  28. RpcContext.getContext().setInvokers((List) selected);
  29. final AtomicInteger count = new AtomicInteger();
  30. final BlockingQueue<Object> ref = new LinkedBlockingQueue<Object>();
  31. // 遍历 selected 列表
  32. for (final Invoker<T> invoker : selected) {
  33. // 为每个 Invoker 创建一个执行线程
  34. executor.execute(new Runnable() {
  35. @Override
  36. public void run() {
  37. try {
  38. // 进行远程调用
  39. Result result = invoker.invoke(invocation);
  40. // 将结果存到阻塞队列中
  41. ref.offer(result);
  42. } catch (Throwable e) {
  43. int value = count.incrementAndGet();
  44. // 仅在 value 大于等于 selected.size() 时,才将异常对象
  45. // 放入阻塞队列中,请大家思考一下为什么要这样做。
  46. if (value >= selected.size()) {
  47. // 将异常对象存入到阻塞队列中
  48. ref.offer(e);
  49. }
  50. }
  51. }
  52. });
  53. }
  54. // ----------------------✨ 分割线2 ✨---------------------- //
  55. try {
  56. // 从阻塞队列中取出远程调用结果
  57. Object ret = ref.poll(timeout, TimeUnit.MILLISECONDS);
  58. // 如果结果类型为 Throwable,则抛出异常
  59. if (ret instanceof Throwable) {
  60. Throwable e = (Throwable) ret;
  61. throw new RpcException(..., "Failed to forking invoke provider ...");
  62. }
  63. // 返回结果
  64. return (Result) ret;
  65. } catch (InterruptedException e) {
  66. throw new RpcException("Failed to forking invoke provider ...");
  67. }
  68. } finally {
  69. RpcContext.getContext().clearAttachments();
  70. }
  71. }
  72. }

ForkingClusterInvoker 的 doInvoker 方法比较长,这里通过两个分割线将整个方法划分为三个逻辑块。从方法开始到分割线1之间的代码主要是用于选出 forks 个 Invoker,为接下来的并发调用提供输入。分割线1和分割线2之间的逻辑通过线程池并发调用多个 Invoker,并将结果存储在阻塞队列中。分割线2到方法结尾之间的逻辑主要用于从阻塞队列中获取返回结果,并对返回结果类型进行判断。如果为异常类型,则直接抛出,否则返回。

以上就是ForkingClusterInvoker 的 doInvoker 方法大致过程。我们在分割线1和分割线2之间的代码上留了一个问题,问题是这样的:为什么要在value >= selected.size()的情况下,才将异常对象添加到阻塞队列中?这里来解答一下。原因是这样的,在并行调用多个服务提供者的情况下,只要有一个服务提供者能够成功返回结果,而其他全部失败。此时 ForkingClusterInvoker 仍应该返回成功的结果,而非抛出异常。在value >= selected.size()时将异常对象放入阻塞队列中,可以保证异常对象不会出现在正常结果的前面,这样可从阻塞队列中优先取出正常的结果。

关于 ForkingClusterInvoker 就先分析到这,接下来分析最后一个 Cluster Invoker。

3.2.6 BroadcastClusterInvoker

本章的最后,我们再来看一下 BroadcastClusterInvoker。BroadcastClusterInvoker 会逐个调用每个服务提供者,如果其中一台报错,在循环调用结束后,BroadcastClusterInvoker 会抛出异常。该类通常用于通知所有提供者更新缓存或日志等本地资源信息。源码如下。

  1. public class BroadcastClusterInvoker<T> extends AbstractClusterInvoker<T> {
  2. @Override
  3. public Result doInvoke(final Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) throws RpcException {
  4. checkInvokers(invokers, invocation);
  5. RpcContext.getContext().setInvokers((List) invokers);
  6. RpcException exception = null;
  7. Result result = null;
  8. // 遍历 Invoker 列表,逐个调用
  9. for (Invoker<T> invoker : invokers) {
  10. try {
  11. // 进行远程调用
  12. result = invoker.invoke(invocation);
  13. } catch (RpcException e) {
  14. exception = e;
  15. logger.warn(e.getMessage(), e);
  16. } catch (Throwable e) {
  17. exception = new RpcException(e.getMessage(), e);
  18. logger.warn(e.getMessage(), e);
  19. }
  20. }
  21. // exception 不为空,则抛出异常
  22. if (exception != null) {
  23. throw exception;
  24. }
  25. return result;
  26. }
  27. }

以上就是 BroadcastClusterInvoker 的代码,比较简单,就不多说了。