From 048fadf2c0fe9fdcb2292b9be975e9f5c0b0ca2c Mon Sep 17 00:00:00 2001 From: "ken.lj" Date: Thu, 27 Feb 2020 10:43:18 +0800 Subject: [PATCH] Distinguish getUrl and getConsumerUrl from Directory. (#5775) * distinguish getUrl and getConsumerUrl from Directory --- .../support/AbstractClusterInvoker.java | 8 +++---- .../support/FailoverClusterInvoker.java | 2 +- .../support/ForkingClusterInvoker.java | 4 ++-- .../support/MergeableClusterInvoker.java | 2 +- .../registry/ZoneAwareClusterInvoker.java | 22 ++++++++++++------- .../support/wrapper/MockClusterInvoker.java | 14 +++++++----- .../support/AbstractClusterInvokerTest.java | 5 +++-- 7 files changed, 34 insertions(+), 23 deletions(-) diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java index 33c036a33b..2a81c7a138 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java @@ -85,11 +85,11 @@ public abstract class AbstractClusterInvoker implements Invoker { @Override public URL getUrl() { - return directory.getUrl(); + return directory.getConsumerUrl(); } - protected URL getConsumerUrl() { - return directory.getConsumerUrl(); + public URL getRegistryUrl() { + return directory.getUrl(); } @Override @@ -255,7 +255,7 @@ public abstract class AbstractClusterInvoker implements Invoker { List> invokers = list(invocation); LoadBalance loadbalance = initLoadBalance(invokers, invocation); - RpcUtils.attachInvocationIdIfAsync(getConsumerUrl(), invocation); + RpcUtils.attachInvocationIdIfAsync(getUrl(), invocation); return doInvoke(invocation, invokers, loadbalance); } diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/FailoverClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/FailoverClusterInvoker.java index e49a99891e..5efe0ce0b9 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/FailoverClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/FailoverClusterInvoker.java @@ -58,7 +58,7 @@ public class FailoverClusterInvoker extends AbstractClusterInvoker { List> copyInvokers = invokers; checkInvokers(copyInvokers, invocation); String methodName = RpcUtils.getMethodName(invocation); - int len = getConsumerUrl().getMethodParameter(methodName, RETRIES_KEY, DEFAULT_RETRIES) + 1; + int len = getUrl().getMethodParameter(methodName, RETRIES_KEY, DEFAULT_RETRIES) + 1; if (len <= 0) { len = 1; } diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/ForkingClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/ForkingClusterInvoker.java index c0a1ddfd4c..7ce7729534 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/ForkingClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/ForkingClusterInvoker.java @@ -65,8 +65,8 @@ public class ForkingClusterInvoker extends AbstractClusterInvoker { try { checkInvokers(invokers, invocation); final List> selected; - final int forks = getConsumerUrl().getParameter(FORKS_KEY, DEFAULT_FORKS); - final int timeout = getConsumerUrl().getParameter(TIMEOUT_KEY, DEFAULT_TIMEOUT); + final int forks = getUrl().getParameter(FORKS_KEY, DEFAULT_FORKS); + final int timeout = getUrl().getParameter(TIMEOUT_KEY, DEFAULT_TIMEOUT); if (forks <= 0 || forks >= invokers.size()) { selected = invokers; } else { diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/MergeableClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/MergeableClusterInvoker.java index 0a9217b233..babc1d52ed 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/MergeableClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/MergeableClusterInvoker.java @@ -63,7 +63,7 @@ public class MergeableClusterInvoker extends AbstractClusterInvoker { @Override protected Result doInvoke(Invocation invocation, List> invokers, LoadBalance loadbalance) throws RpcException { checkInvokers(invokers, invocation); - String merger = getConsumerUrl().getMethodParameter(invocation.getMethodName(), MERGER_KEY); + String merger = getUrl().getMethodParameter(invocation.getMethodName(), MERGER_KEY); if (ConfigUtils.isEmpty(merger)) { // If a method doesn't have a merger, only invoke one Group for (final Invoker invoker : invokers) { if (invoker.isAvailable()) { diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/registry/ZoneAwareClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/registry/ZoneAwareClusterInvoker.java index 09c95c25a4..f9f2ae72ea 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/registry/ZoneAwareClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/registry/ZoneAwareClusterInvoker.java @@ -26,6 +26,7 @@ import org.apache.dubbo.rpc.RpcException; import org.apache.dubbo.rpc.cluster.Directory; import org.apache.dubbo.rpc.cluster.LoadBalance; import org.apache.dubbo.rpc.cluster.support.AbstractClusterInvoker; +import org.apache.dubbo.rpc.cluster.support.wrapper.MockClusterInvoker; import java.util.List; import java.util.stream.Collectors; @@ -58,24 +59,28 @@ public class ZoneAwareClusterInvoker extends AbstractClusterInvoker { public Result doInvoke(Invocation invocation, final List> invokers, LoadBalance loadbalance) throws RpcException { // First, pick the invoker (XXXClusterInvoker) that comes from the local registry, distinguish by a 'preferred' key. for (Invoker invoker : invokers) { - if (invoker.isAvailable() && invoker.getUrl().getParameter(REGISTRY_KEY + "." + PREFERRED_KEY, false)) { - return invoker.invoke(invocation); + // FIXME, the invoker is a cluster invoker representing one Registry, so it will automatically wrapped by MockClusterInvoker. + MockClusterInvoker mockClusterInvoker = (MockClusterInvoker) invoker; + if (mockClusterInvoker.isAvailable() && mockClusterInvoker.getRegistryUrl() + .getParameter(REGISTRY_KEY + "." + PREFERRED_KEY, false)) { + return mockClusterInvoker.invoke(invocation); } } - // providers in the registry with the same + // providers in the registry with the same zone String zone = (String) invocation.getAttachment(REGISTRY_ZONE); if (StringUtils.isNotEmpty(zone)) { for (Invoker invoker : invokers) { - if (invoker.isAvailable() && zone.equals(invoker.getUrl().getParameter(REGISTRY_KEY + "." + ZONE_KEY))) { - return invoker.invoke(invocation); + MockClusterInvoker mockClusterInvoker = (MockClusterInvoker) invoker; + if (mockClusterInvoker.isAvailable() && zone.equals(mockClusterInvoker.getRegistryUrl().getParameter(REGISTRY_KEY + "." + ZONE_KEY))) { + return mockClusterInvoker.invoke(invocation); } } String force = (String) invocation.getAttachment(REGISTRY_ZONE_FORCE); if (StringUtils.isNotEmpty(force) && "true".equalsIgnoreCase(force)) { throw new IllegalStateException("No registry instance in zone or no available providers in the registry, zone: " + zone - + ", registries: " + invokers.stream().map(i -> i.getUrl().toString()).collect(Collectors.joining(","))); + + ", registries: " + invokers.stream().map(invoker -> ((MockClusterInvoker) invoker).getRegistryUrl().toString()).collect(Collectors.joining(","))); } } @@ -88,8 +93,9 @@ public class ZoneAwareClusterInvoker extends AbstractClusterInvoker { // If none of the invokers has a preferred signal or is picked by the loadbalancer, pick the first one available. for (Invoker invoker : invokers) { - if (invoker.isAvailable()) { - return invoker.invoke(invocation); + MockClusterInvoker mockClusterInvoker = (MockClusterInvoker) invoker; + if (mockClusterInvoker.isAvailable()) { + return mockClusterInvoker.invoke(invocation); } } throw new RpcException("No provider available in " + invokers); diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/wrapper/MockClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/wrapper/MockClusterInvoker.java index 6b12e1c2c3..afdeddee9e 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/wrapper/MockClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/wrapper/MockClusterInvoker.java @@ -50,6 +50,10 @@ public class MockClusterInvoker implements Invoker { @Override public URL getUrl() { + return directory.getConsumerUrl(); + } + + public URL getRegistryUrl() { return directory.getUrl(); } @@ -72,13 +76,13 @@ public class MockClusterInvoker implements Invoker { public Result invoke(Invocation invocation) throws RpcException { Result result = null; - String value = directory.getConsumerUrl().getMethodParameter(invocation.getMethodName(), MOCK_KEY, Boolean.FALSE.toString()).trim(); + String value = getUrl().getMethodParameter(invocation.getMethodName(), MOCK_KEY, Boolean.FALSE.toString()).trim(); if (value.length() == 0 || "false".equalsIgnoreCase(value)) { //no mock result = this.invoker.invoke(invocation); } else if (value.startsWith("force")) { if (logger.isWarnEnabled()) { - logger.warn("force-mock: " + invocation.getMethodName() + " force-mock enabled , url : " + directory.getConsumerUrl()); + logger.warn("force-mock: " + invocation.getMethodName() + " force-mock enabled , url : " + getUrl()); } //force:direct mock result = doMockInvoke(invocation, null); @@ -103,7 +107,7 @@ public class MockClusterInvoker implements Invoker { } if (logger.isWarnEnabled()) { - logger.warn("fail-mock: " + invocation.getMethodName() + " fail-mock enabled , url : " + directory.getConsumerUrl(), e); + logger.warn("fail-mock: " + invocation.getMethodName() + " fail-mock enabled , url : " + getUrl(), e); } result = doMockInvoke(invocation, e); } @@ -118,7 +122,7 @@ public class MockClusterInvoker implements Invoker { List> mockInvokers = selectMockInvoker(invocation); if (CollectionUtils.isEmpty(mockInvokers)) { - minvoker = (Invoker) new MockInvoker(directory.getConsumerUrl(), directory.getInterface()); + minvoker = (Invoker) new MockInvoker(getUrl(), directory.getInterface()); } else { minvoker = mockInvokers.get(0); } @@ -165,7 +169,7 @@ public class MockClusterInvoker implements Invoker { } catch (RpcException e) { if (logger.isInfoEnabled()) { logger.info("Exception when try to invoke mock. Get mock invokers error for service:" - + directory.getConsumerUrl().getServiceInterface() + ", method:" + invocation.getMethodName() + + getUrl().getServiceInterface() + ", method:" + invocation.getMethodName() + ", will construct a new mock with 'new MockInvoker()'.", e); } } diff --git a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java index e94ebabeb5..955ec0ef5b 100644 --- a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java +++ b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java @@ -19,6 +19,7 @@ package org.apache.dubbo.rpc.cluster.support; import org.apache.dubbo.common.URL; import org.apache.dubbo.common.extension.ExtensionLoader; import org.apache.dubbo.common.utils.NetUtils; +import org.apache.dubbo.common.utils.StringUtils; import org.apache.dubbo.rpc.Invocation; import org.apache.dubbo.rpc.Invoker; import org.apache.dubbo.rpc.Result; @@ -65,7 +66,6 @@ public class AbstractClusterInvokerTest { StaticDirectory dic; RpcInvocation invocation = new RpcInvocation(); URL url = URL.valueOf("registry://localhost:9090/org.apache.dubbo.rpc.cluster.support.AbstractClusterInvokerTest.IHelloService?refer=" + URL.encode("application=abstractClusterInvokerTest")); - URL tmpUrl = url.removeParameter(REFER_KEY).removeParameter(MONITOR_KEY); Invoker invoker1; Invoker invoker2; @@ -124,7 +124,6 @@ public class AbstractClusterInvokerTest { invokers.add(invoker1); dic = new StaticDirectory(url, invokers, null); - cluster = new AbstractClusterInvoker(dic) { @Override protected Result doInvoke(Invocation invocation, List invokers, LoadBalance loadbalance) @@ -226,6 +225,8 @@ public class AbstractClusterInvokerTest { @Test public void testCloseAvailablecheck() { LoadBalance lb = mock(LoadBalance.class); + Map queryMap = StringUtils.parseQueryString(url.getParameterAndDecoded(REFER_KEY)); + URL tmpUrl = url.addParameters(queryMap).removeParameter(MONITOR_KEY); given(lb.select(invokers, tmpUrl, invocation)).willReturn(invoker1); initlistsize5();