diff --git a/dubbo-dependencies-bom/pom.xml b/dubbo-dependencies-bom/pom.xml
index 089a617797..0bb37e2663 100644
--- a/dubbo-dependencies-bom/pom.xml
+++ b/dubbo-dependencies-bom/pom.xml
@@ -143,7 +143,7 @@
8.5.78
0.5.3
2.1.2
- 1.47.0
+ 1.41.0
0.8.1
1.2.1
diff --git a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/PilotExchanger.java b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/PilotExchanger.java
index dbb30e7586..ce84aab722 100644
--- a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/PilotExchanger.java
+++ b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/PilotExchanger.java
@@ -26,85 +26,84 @@ import org.apache.dubbo.registry.xds.util.protocol.message.Endpoint;
import org.apache.dubbo.registry.xds.util.protocol.message.EndpointResult;
import org.apache.dubbo.registry.xds.util.protocol.message.ListenerResult;
import org.apache.dubbo.registry.xds.util.protocol.message.RouteResult;
+import org.apache.dubbo.rpc.model.ApplicationModel;
import java.util.Collections;
+import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Consumer;
public class PilotExchanger {
- private final XdsChannel xdsChannel;
+ protected final XdsChannel xdsChannel;
- private final RdsProtocol rdsProtocol;
+ protected final LdsProtocol ldsProtocol;
- private final EdsProtocol edsProtocol;
+ protected final RdsProtocol rdsProtocol;
- private ListenerResult listenerResult;
+ protected final EdsProtocol edsProtocol;
- private RouteResult routeResult;
+ protected Map listenerResult;
- private final AtomicLong observeRouteRequest = new AtomicLong(-1);
+ protected Map routeResult;
- private final Map domainObserveRequest = new ConcurrentHashMap<>();
+ private final AtomicBoolean isRdsObserve = new AtomicBoolean(false);
+ private final HashSet domainObserveRequest = new HashSet<>();
private final Map>>> domainObserveConsumer = new ConcurrentHashMap<>();
- private PilotExchanger(URL url) {
+ protected PilotExchanger(URL url) {
xdsChannel = new XdsChannel(url);
- int pollingPoolSize = url.getParameter("pollingPoolSize", 10);
int pollingTimeout = url.getParameter("pollingTimeout", 10);
- LdsProtocol ldsProtocol = new LdsProtocol(xdsChannel, NodeBuilder.build(), pollingPoolSize, pollingTimeout);
- this.rdsProtocol = new RdsProtocol(xdsChannel, NodeBuilder.build(), pollingPoolSize, pollingTimeout);
- this.edsProtocol = new EdsProtocol(xdsChannel, NodeBuilder.build(), pollingPoolSize, pollingTimeout);
+ ApplicationModel applicationModel = url.getOrDefaultApplicationModel();
+ this.ldsProtocol = new LdsProtocol(xdsChannel, NodeBuilder.build(), pollingTimeout, applicationModel);
+ this.rdsProtocol = new RdsProtocol(xdsChannel, NodeBuilder.build(), pollingTimeout, applicationModel);
+ this.edsProtocol = new EdsProtocol(xdsChannel, NodeBuilder.build(), pollingTimeout, applicationModel);
this.listenerResult = ldsProtocol.getListeners();
- this.routeResult = rdsProtocol.getResource(listenerResult.getRouteConfigNames());
+ this.routeResult = rdsProtocol.getResource(listenerResult.values().iterator().next().getRouteConfigNames());
// Observer RDS update
- if (CollectionUtils.isNotEmpty(listenerResult.getRouteConfigNames())) {
- this.observeRouteRequest.set(createRouteObserve());
+ if (CollectionUtils.isNotEmpty(listenerResult.values().iterator().next().getRouteConfigNames())) {
+ createRouteObserve();
+ isRdsObserve.set(true);
}
+
// Observe LDS updated
ldsProtocol.observeListeners((newListener) -> {
// update local cache
if (!newListener.equals(listenerResult)) {
this.listenerResult = newListener;
// update RDS observation
- synchronized (observeRouteRequest) {
- if (observeRouteRequest.get() == -1) {
- this.observeRouteRequest.set(createRouteObserve());
- } else {
- rdsProtocol.updateObserve(observeRouteRequest.get(), newListener.getRouteConfigNames());
- }
+ if (isRdsObserve.get()) {
+ createRouteObserve();
}
}
});
}
- private long createRouteObserve() {
- return rdsProtocol.observeResource(listenerResult.getRouteConfigNames(), (newResult) -> {
+ private void createRouteObserve() {
+ rdsProtocol.observeResource(listenerResult.values().iterator().next().getRouteConfigNames(), (newResult) -> {
// check if observed domain update ( will update endpoint observation )
domainObserveConsumer.forEach((domain, consumer) -> {
- Set newRoute = newResult.searchDomain(domain);
- if (!routeResult.searchDomain(domain).equals(newRoute)) {
- // routers in observed domain has been updated
- Long domainRequest = domainObserveRequest.get(domain);
- if (domainRequest == null) {
- // router list is empty when observeEndpoints() called and domainRequest has not been created yet
- // create new observation
- doObserveEndpoints(domain);
- } else {
- // update observation by domainRequest
- edsProtocol.updateObserve(domainRequest, newRoute);
+ newResult.values().forEach(o -> {
+ Set newRoute = o.searchDomain(domain);
+ for (Map.Entry entry: routeResult.entrySet()) {
+ if (!entry.getValue().searchDomain(domain).equals(newRoute)) {
+ // routers in observed domain has been updated
+// Long domainRequest = domainObserveRequest.get(domain);
+ // router list is empty when observeEndpoints() called and domainRequest has not been created yet
+ // create new observation
+ doObserveEndpoints(domain);
+ }
}
- }
+ });
});
- // update local cache
routeResult = newResult;
- });
+ }, false);
}
public static PilotExchanger initialize(URL url) {
@@ -116,17 +115,25 @@ public class PilotExchanger {
}
public Set getServices() {
- return routeResult.getDomains();
+ Set domains = new HashSet<>();
+ for (Map.Entry entry: routeResult.entrySet()) {
+ domains.addAll(entry.getValue().getDomains());
+ }
+ return domains;
}
public Set getEndpoints(String domain) {
- Set cluster = routeResult.searchDomain(domain);
- if (CollectionUtils.isNotEmpty(cluster)) {
- EndpointResult endpoint = edsProtocol.getResource(cluster);
- return endpoint.getEndpoints();
- } else {
- return Collections.emptySet();
+ Set endpoints = new HashSet<>();
+ for (Map.Entry entry: routeResult.entrySet()) {
+ Set cluster = entry.getValue().searchDomain(domain);
+ if (CollectionUtils.isNotEmpty(cluster)) {
+ Map endpointResultList = edsProtocol.getResource(cluster);
+ endpointResultList.forEach((k, v) -> endpoints.addAll(v.getEndpoints()));
+ } else {
+ return Collections.emptySet();
+ }
}
+ return endpoints;
}
public void observeEndpoints(String domain, Consumer> consumer) {
@@ -139,24 +146,29 @@ public class PilotExchanger {
v.add(consumer);
return v;
});
- if (!domainObserveRequest.containsKey(domain)) {
+ if (!domainObserveRequest.contains(domain)) {
doObserveEndpoints(domain);
}
}
private void doObserveEndpoints(String domain) {
- Set router = routeResult.searchDomain(domain);
- // if router is empty, do nothing
- // observation will be created when RDS updates
- if (CollectionUtils.isNotEmpty(router)) {
- long endpointRequest =
+ for (Map.Entry entry: routeResult.entrySet()) {
+ Set router = entry.getValue().searchDomain(domain);
+ // if router is empty, do nothing
+ // observation will be created when RDS updates
+ if (CollectionUtils.isNotEmpty(router)) {
edsProtocol.observeResource(
router,
- endpointResult ->
- // notify consumers
- domainObserveConsumer.get(domain).forEach(
- consumer1 -> consumer1.accept(endpointResult.getEndpoints())));
- domainObserveRequest.put(domain, endpointRequest);
+ (endpointResultMap) -> {
+ endpointResultMap.forEach((k, v) -> {
+ // notify consumers
+ domainObserveConsumer.get(domain).forEach(
+ consumer1 -> consumer1.accept(v.getEndpoints()));
+ });
+ }, false);
+ domainObserveRequest.add(domain);
+ }
}
+
}
}
diff --git a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/XdsChannel.java b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/XdsChannel.java
index e58eed0172..d4e069a364 100644
--- a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/XdsChannel.java
+++ b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/XdsChannel.java
@@ -49,14 +49,25 @@ public class XdsChannel {
private static final String USE_AGENT = "use-agent";
+ private URL url;
+
private static final String SECURE = "secure";
private static final String PLAINTEXT = "plaintext";
private final ManagedChannel channel;
- protected XdsChannel(URL url) {
+ public URL getUrl() {
+ return url;
+ }
+
+ public ManagedChannel getChannel() {
+ return channel;
+ }
+
+ public XdsChannel(URL url) {
ManagedChannel managedChannel = null;
+ this.url = url;
try {
if (!url.getParameter(USE_AGENT, false)) {
if(PLAINTEXT.equals(url.getParameter(SECURE))){
diff --git a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
index fdb98c3410..6fcd422635 100644
--- a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
+++ b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
@@ -16,76 +16,77 @@
*/
package org.apache.dubbo.registry.xds.util.protocol;
-import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
-import org.apache.dubbo.common.logger.LoggerFactory;
-import org.apache.dubbo.common.utils.NamedThreadFactory;
-import org.apache.dubbo.registry.xds.util.XdsChannel;
import io.envoyproxy.envoy.config.core.v3.Node;
import io.envoyproxy.envoy.service.discovery.v3.DiscoveryRequest;
import io.envoyproxy.envoy.service.discovery.v3.DiscoveryResponse;
import io.grpc.stub.StreamObserver;
+import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
+import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.threadpool.manager.FrameworkExecutorRepository;
+import org.apache.dubbo.registry.xds.util.XdsChannel;
+import org.apache.dubbo.rpc.model.ApplicationModel;
+import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutionException;
import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ScheduledFuture;
-import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.locks.ReentrantLock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Consumer;
+import java.util.stream.Collectors;
import static org.apache.dubbo.common.constants.LoggerCodeConstants.REGISTRY_ERROR_REQUEST_XDS;
+import static org.apache.dubbo.common.constants.LoggerCodeConstants.PROTOCOL_FAILED_REQUEST;
+import static org.apache.dubbo.common.constants.LoggerCodeConstants.INTERNAL_INTERRUPTED;
public abstract class AbstractProtocol> implements XdsProtocol {
private static final ErrorTypeAwareLogger logger = LoggerFactory.getErrorTypeAwareLogger(AbstractProtocol.class);
- protected final XdsChannel xdsChannel;
+ protected XdsChannel xdsChannel;
protected final Node node;
- /**
- * Store Request Parameter ( resourceNames )
- * K - requestId, V - resourceNames
- */
- protected final Map> requestParam = new ConcurrentHashMap<>();
+ private final int checkInterval;
- /**
- * Store ADS Request Observer ( StreamObserver in Streaming Request )
- * K - requestId, V - StreamObserver
- */
- private final Map> requestObserverMap = new ConcurrentHashMap<>();
+ protected final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
- /**
- * Store Delta-ADS Request Observer ( StreamObserver in Streaming Request )
- * K - requestId, V - StreamObserver
- */
- private final Map> observeScheduledMap = new ConcurrentHashMap<>();
+ protected final ReentrantReadWriteLock.ReadLock readLock = lock.readLock();
- /**
- * Store CompletableFuture for Request ( used to fetch async result in ResponseObserver )
- * K - requestId, V - CompletableFuture
- */
- private final Map> streamResult = new ConcurrentHashMap<>();
+ protected final ReentrantReadWriteLock.WriteLock writeLock = lock.writeLock();
- private final ScheduledExecutorService pollingExecutor;
+ protected Set observeResourcesName;
- private final int pollingTimeout;
+ public static final String emptyResourceName = "emptyResourcesName";
+ private final ReentrantLock resourceLock = new ReentrantLock();
- protected final static AtomicLong requestId = new AtomicLong(0);
+ protected Map, List>>> consumerObserveMap = new ConcurrentHashMap<>();
- public AbstractProtocol(XdsChannel xdsChannel, Node node, int pollingPoolSize, int pollingTimeout) {
+ public Map, List>>> getConsumerObserveMap() {
+ return consumerObserveMap;
+ }
+ private final ApplicationModel applicationModel;
+
+ public AbstractProtocol(XdsChannel xdsChannel, Node node, int checkInterval, ApplicationModel applicationModel) {
this.xdsChannel = xdsChannel;
this.node = node;
- this.pollingExecutor = new ScheduledThreadPoolExecutor(pollingPoolSize, new NamedThreadFactory("Dubbo-registry-xds"));
- this.pollingTimeout = pollingTimeout;
+ this.checkInterval = checkInterval;
+ this.applicationModel = applicationModel;
}
+ protected Map resourcesMap = new ConcurrentHashMap<>();
+
+ protected StreamObserver requestObserver;
+
/**
* Abstract method to obtain Type-URL from sub-class
*
@@ -93,103 +94,126 @@ public abstract class AbstractProtocol> implements
*/
public abstract String getTypeUrl();
+ public boolean isCacheExistResource(Set resourceNames) {
+ for (String resourceName : resourceNames) {
+ if ("".equals(resourceName)) {
+ continue;
+ }
+ if (!resourcesMap.containsKey(resourceName)) {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ public T getCacheResource(String resourceName) {
+ if (resourceName == null || resourceName.length() == 0) {
+ return null;
+ }
+ return resourcesMap.get(resourceName);
+ }
+
+
@Override
- public T getResource(Set resourceNames) {
- long request = requestId.getAndIncrement();
+ public Map getResource(Set resourceNames) {
resourceNames = resourceNames == null ? Collections.emptySet() : resourceNames;
- // Store Request Parameter, which will be used for ACK
- requestParam.put(request, resourceNames);
-
- // create observer
- StreamObserver requestObserver = xdsChannel.createDeltaDiscoveryRequest(new ResponseObserver(request));
-
- // use future to get async result
- CompletableFuture future = new CompletableFuture<>();
- requestObserverMap.put(request, requestObserver);
- streamResult.put(request, future);
-
- // send request to control panel
- requestObserver.onNext(buildDiscoveryRequest(resourceNames));
-
- try {
- // get result
- return future.get();
- } catch (InterruptedException | ExecutionException e) {
- logger.error(REGISTRY_ERROR_REQUEST_XDS, "", "", "Error occur when request control panel.");
- return null;
- } finally {
- // close observer
- //requestObserver.onCompleted();
-
- // remove temp
- streamResult.remove(request);
- requestObserverMap.remove(request);
- requestParam.remove(request);
+ if (!resourceNames.isEmpty() && isCacheExistResource(resourceNames)) {
+ return getResourceFromCache(resourceNames);
+ } else {
+ return getResourceFromRemote(resourceNames);
}
}
- @Override
- public long observeResource(Set resourceNames, Consumer consumer) {
- long request = requestId.getAndIncrement();
- resourceNames = resourceNames == null ? Collections.emptySet() : resourceNames;
-
- // Store Request Parameter, which will be used for ACK
- requestParam.put(request, resourceNames);
-
- // call once for full data
- consumer.accept(getResource(resourceNames));
-
- // channel reused
- StreamObserver requestObserver = xdsChannel.createDeltaDiscoveryRequest(new ResponseObserver(request));
- requestObserverMap.put(request, requestObserver);
-
- ScheduledFuture> scheduledFuture = pollingExecutor.scheduleAtFixedRate(() -> {
- try {
- // origin request, may changed by updateObserve
- Set names = requestParam.get(request);
-
- // use future to get async result, future complete on StreamObserver onNext
- CompletableFuture future = new CompletableFuture<>();
- streamResult.put(request, future);
-
- // observer reused
- StreamObserver observer = requestObserverMap.get(request);
-
- if (observer == null) {
- observer = xdsChannel.createDeltaDiscoveryRequest(new ResponseObserver(request));
- requestObserverMap.put(request, observer);
- }
-
- // send request to control panel
- observer.onNext(buildDiscoveryRequest(names));
-
- try {
- // get result
- consumer.accept(future.get());
- } catch (InterruptedException | ExecutionException e) {
- logger.error(REGISTRY_ERROR_REQUEST_XDS, "", "", "Error occur when request control panel.");
- } finally {
- // close observer
- //requestObserver.onCompleted();
-
- // remove temp
- streamResult.remove(request);
- }
- } catch (Throwable t) {
- logger.error(REGISTRY_ERROR_REQUEST_XDS, "", "", "Error when requesting observe data. Type: " + getTypeUrl(), t);
- }
- }, pollingTimeout, pollingTimeout, TimeUnit.SECONDS);
-
- observeScheduledMap.put(request, scheduledFuture);
-
- return request;
+ private Map getResourceFromCache(Set resourceNames) {
+ return resourceNames.stream()
+ .collect(Collectors.toMap(k -> k, this::getCacheResource));
}
- @Override
- public void updateObserve(long request, Set resourceNames) {
- // send difference in resourceNames
- requestParam.put(request, resourceNames);
+ public Map getResourceFromRemote(Set resourceNames) {
+ try {
+ resourceLock.lock();
+ CompletableFuture