This is an automated email from the ASF dual-hosted git repository.
liujun pushed a commit to branch 3.0
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.0 by this push:
new ea0506a Collecting registration and consumption status inside process
(#7714)
ea0506a is described below
commit ea0506aa65f230ba805b482370eee6a871f04b13
Author: ken.lj <[email protected]>
AuthorDate: Mon May 10 13:56:04 2021 +0800
Collecting registration and consumption status inside process (#7714)
---
dubbo-common/pom.xml | 4 ++
.../status/reporter/FrameworkStatusReporter.java | 79 ++++++++++++++++++++++
.../dubbo/config/utils/ConfigValidationUtils.java | 10 ++-
.../DynamicConfigurationServiceNameMapping.java | 2 -
.../apache/dubbo/metadata/ServiceNameMapping.java | 2 +-
.../dubbo/metadata/WritableMetadataService.java | 3 +
.../client/AbstractServiceDiscoveryFactory.java | 9 ++-
.../registry/client/DefaultServiceInstance.java | 10 +++
.../client/ServiceDiscoveryRegistryDirectory.java | 2 +-
.../store/InMemoryWritableMetadataService.java | 5 ++
.../DefaultMigrationAddressComparator.java | 38 +++++++++--
.../migration/MigrationAddressComparator.java | 3 +
.../client/migration/MigrationInvoker.java | 55 +++++++++------
.../registry/integration/RegistryDirectory.java | 6 +-
14 files changed, 191 insertions(+), 37 deletions(-)
diff --git a/dubbo-common/pom.xml b/dubbo-common/pom.xml
index b424661..f32317c 100644
--- a/dubbo-common/pom.xml
+++ b/dubbo-common/pom.xml
@@ -64,6 +64,10 @@
<artifactId>fastjson</artifactId>
</dependency>
<dependency>
+ <groupId>com.google.code.gson</groupId>
+ <artifactId>gson</artifactId>
+ </dependency>
+ <dependency>
<groupId>commons-io</groupId>
<artifactId>commons-io</artifactId>
</dependency>
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/status/reporter/FrameworkStatusReporter.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/status/reporter/FrameworkStatusReporter.java
new file mode 100644
index 0000000..6d32e9f
--- /dev/null
+++
b/dubbo-common/src/main/java/org/apache/dubbo/common/status/reporter/FrameworkStatusReporter.java
@@ -0,0 +1,79 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.common.status.reporter;
+
+import org.apache.dubbo.common.extension.ExtensionLoader;
+import org.apache.dubbo.common.extension.SPI;
+import org.apache.dubbo.common.logger.Logger;
+import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.utils.CollectionUtils;
+import org.apache.dubbo.rpc.model.ApplicationModel;
+
+import com.google.gson.Gson;
+
+import java.util.HashMap;
+import java.util.Set;
+
+@SPI
+public interface FrameworkStatusReporter {
+ static final Gson gson = new Gson();
+ Logger logger = LoggerFactory.getLogger(FrameworkStatusReporter.class);
+ String REGISTRATION_STATUS = "registration";
+ String ADDRESS_CONSUMPTION_STATUS = "consumption";
+
+ void report(String type, Object obj);
+
+ static void reportRegistrationStatus(Object obj) {
+ doReport(REGISTRATION_STATUS, obj);
+ }
+
+ static void reportConsumptionStatus(Object obj) {
+ doReport(ADDRESS_CONSUMPTION_STATUS, obj);
+ }
+
+ static void doReport(String type, Object obj) {
+ // TODO, report asynchronously
+ try {
+ Set<FrameworkStatusReporter> reporters =
ExtensionLoader.getExtensionLoader(FrameworkStatusReporter.class).getSupportedExtensionInstances();
+ if (CollectionUtils.isNotEmpty(reporters)) {
+ FrameworkStatusReporter reporter = reporters.iterator().next();
+ reporter.report(type, obj);
+ }
+ } catch (Exception e) {
+ logger.info("Report " + type + " status failed because of " +
e.getMessage());
+ }
+ }
+
+ static String createRegistrationReport(String status) {
+ return "{\"application\":\"" +
+ ApplicationModel.getName() +
+ "\",\"status\":\"" +
+ status +
+ "\"}";
+ }
+
+ static String createConsumptionReport(String interfaceName, String
version, String group, String status) {
+ HashMap<String, String> migrationStatus = new HashMap<>();
+ migrationStatus.put("type", "consumption");
+ migrationStatus.put("application", ApplicationModel.getName());
+ migrationStatus.put("service", interfaceName);
+ migrationStatus.put("version", version);
+ migrationStatus.put("group", group);
+ migrationStatus.put("status", status);
+ return gson.toJson(migrationStatus);
+ }
+}
diff --git
a/dubbo-config/dubbo-config-api/src/main/java/org/apache/dubbo/config/utils/ConfigValidationUtils.java
b/dubbo-config/dubbo-config-api/src/main/java/org/apache/dubbo/config/utils/ConfigValidationUtils.java
index 6919652..2f68774 100644
---
a/dubbo-config/dubbo-config-api/src/main/java/org/apache/dubbo/config/utils/ConfigValidationUtils.java
+++
b/dubbo-config/dubbo-config-api/src/main/java/org/apache/dubbo/config/utils/ConfigValidationUtils.java
@@ -24,6 +24,7 @@ import org.apache.dubbo.common.logger.Logger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.serialize.Serialization;
import org.apache.dubbo.common.status.StatusChecker;
+import org.apache.dubbo.common.status.reporter.FrameworkStatusReporter;
import org.apache.dubbo.common.threadpool.ThreadPool;
import org.apache.dubbo.common.utils.CollectionUtils;
import org.apache.dubbo.common.utils.ConfigUtils;
@@ -98,6 +99,7 @@ import static
org.apache.dubbo.common.constants.RegistryConstants.REGISTRY_TYPE_
import static
org.apache.dubbo.common.constants.RegistryConstants.SERVICE_REGISTRY_PROTOCOL;
import static org.apache.dubbo.common.constants.RemotingConstants.BACKUP_KEY;
import static
org.apache.dubbo.common.extension.ExtensionLoader.getExtensionLoader;
+import static
org.apache.dubbo.common.status.reporter.FrameworkStatusReporter.createRegistrationReport;
import static org.apache.dubbo.config.Constants.ARCHITECTURE;
import static org.apache.dubbo.config.Constants.CONTEXTPATH_KEY;
import static org.apache.dubbo.config.Constants.DUBBO_IP_TO_REGISTRY;
@@ -216,8 +218,9 @@ public class ConfigValidationUtils {
registryList.forEach(registryURL -> {
if (provider) {
// for registries enabled service discovery, automatically
register interface compatible addresses.
+ String registerMode;
if
(SERVICE_REGISTRY_PROTOCOL.equals(registryURL.getProtocol())) {
- String registerMode =
registryURL.getParameter(REGISTER_MODE_KEY,
ConfigurationUtils.getDynamicGlobalConfiguration().getString(DUBBO_REGISTER_MODE_DEFAULT_KEY,
DEFAULT_REGISTER_MODE_INSTANCE));
+ registerMode = registryURL.getParameter(REGISTER_MODE_KEY,
ConfigurationUtils.getDynamicGlobalConfiguration().getString(DUBBO_REGISTER_MODE_DEFAULT_KEY,
DEFAULT_REGISTER_MODE_INSTANCE));
if (!isValidRegisterMode(registerMode)) {
registerMode = DEFAULT_REGISTER_MODE_INSTANCE;
}
@@ -231,7 +234,7 @@ public class ConfigValidationUtils {
result.add(interfaceCompatibleRegistryURL);
}
} else {
- String registerMode =
registryURL.getParameter(REGISTER_MODE_KEY,
ConfigurationUtils.getDynamicGlobalConfiguration().getString(DUBBO_REGISTER_MODE_DEFAULT_KEY,
DEFAULT_REGISTER_MODE_INTERFACE));
+ registerMode = registryURL.getParameter(REGISTER_MODE_KEY,
ConfigurationUtils.getDynamicGlobalConfiguration().getString(DUBBO_REGISTER_MODE_DEFAULT_KEY,
DEFAULT_REGISTER_MODE_INTERFACE));
if (!isValidRegisterMode(registerMode)) {
registerMode = DEFAULT_REGISTER_MODE_INTERFACE;
}
@@ -248,10 +251,13 @@ public class ConfigValidationUtils {
result.add(registryURL);
}
}
+
+
FrameworkStatusReporter.reportRegistrationStatus(createRegistrationReport(registerMode));
} else {
result.add(registryURL);
}
});
+
return result;
}
diff --git
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/DynamicConfigurationServiceNameMapping.java
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/DynamicConfigurationServiceNameMapping.java
index bac064a..f524fcf 100644
---
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/DynamicConfigurationServiceNameMapping.java
+++
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/DynamicConfigurationServiceNameMapping.java
@@ -40,8 +40,6 @@ import static
org.apache.dubbo.rpc.model.ApplicationModel.getName;
*/
public class DynamicConfigurationServiceNameMapping implements
ServiceNameMapping {
- public static String DEFAULT_MAPPING_GROUP = "mapping";
-
private static final List<String> IGNORED_SERVICE_INTERFACES =
asList(MetadataService.class.getName());
private final Logger logger = LoggerFactory.getLogger(getClass());
diff --git
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/ServiceNameMapping.java
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/ServiceNameMapping.java
index a233910..288712e 100644
---
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/ServiceNameMapping.java
+++
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/ServiceNameMapping.java
@@ -27,7 +27,6 @@ import static
org.apache.dubbo.common.constants.CommonConstants.DUBBO;
import static org.apache.dubbo.common.constants.CommonConstants.PROTOCOL_KEY;
import static
org.apache.dubbo.common.extension.ExtensionLoader.getExtensionLoader;
import static org.apache.dubbo.common.utils.StringUtils.SLASH;
-import static
org.apache.dubbo.metadata.DynamicConfigurationServiceNameMapping.DEFAULT_MAPPING_GROUP;
/**
* The interface for Dubbo service name Mapping
@@ -36,6 +35,7 @@ import static
org.apache.dubbo.metadata.DynamicConfigurationServiceNameMapping.D
*/
@SPI("config")
public interface ServiceNameMapping {
+ String DEFAULT_MAPPING_GROUP = "mapping";
/**
* Map the specified Dubbo service interface, group, version and protocol
to current Dubbo service name
diff --git
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/WritableMetadataService.java
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/WritableMetadataService.java
index 1782ef7..222da5f 100644
---
a/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/WritableMetadataService.java
+++
b/dubbo-metadata/dubbo-metadata-api/src/main/java/org/apache/dubbo/metadata/WritableMetadataService.java
@@ -21,6 +21,7 @@ import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.extension.SPI;
import org.apache.dubbo.rpc.model.ApplicationModel;
+import java.util.Map;
import java.util.Set;
import static
org.apache.dubbo.common.extension.ExtensionLoader.getExtensionLoader;
@@ -93,6 +94,8 @@ public interface WritableMetadataService extends
MetadataService {
Set<String> removeCachedMapping(String serviceKey);
+ Map<String, Set<String>> getCachedMapping();
+
MetadataInfo getDefaultMetadataInfo();
/**
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/AbstractServiceDiscoveryFactory.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/AbstractServiceDiscoveryFactory.java
index 1088b7b..90d8127 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/AbstractServiceDiscoveryFactory.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/AbstractServiceDiscoveryFactory.java
@@ -18,6 +18,9 @@ package org.apache.dubbo.registry.client;
import org.apache.dubbo.common.URL;
+import java.util.Collections;
+import java.util.LinkedList;
+import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -30,7 +33,11 @@ import java.util.concurrent.ConcurrentMap;
*/
public abstract class AbstractServiceDiscoveryFactory implements
ServiceDiscoveryFactory {
- private final ConcurrentMap<String, ServiceDiscovery> discoveries = new
ConcurrentHashMap<>();
+ private static final ConcurrentMap<String, ServiceDiscovery> discoveries =
new ConcurrentHashMap<>();
+
+ public static List<ServiceDiscovery> getAllServiceDiscoveries() {
+ return Collections.unmodifiableList(new
LinkedList<>(discoveries.values()));
+ }
@Override
public ServiceDiscovery getServiceDiscovery(URL registryURL) {
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/DefaultServiceInstance.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/DefaultServiceInstance.java
index 1bcf284..10fe957 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/DefaultServiceInstance.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/DefaultServiceInstance.java
@@ -40,6 +40,8 @@ public class DefaultServiceInstance implements
ServiceInstance {
private static final long serialVersionUID = 1149677083747278100L;
+ private String rawAddress;
+
private String serviceName;
private String host;
@@ -87,6 +89,10 @@ public class DefaultServiceInstance implements
ServiceInstance {
this.healthy = true;
}
+ public void setRawAddress(String rawAddress) {
+ this.rawAddress = rawAddress;
+ }
+
public DefaultServiceInstance(String serviceName) {
this.serviceName = serviceName;
}
@@ -249,6 +255,10 @@ public class DefaultServiceInstance implements
ServiceInstance {
@Override
public String toString() {
+ return rawAddress == null ? toFullString() : rawAddress;
+ }
+
+ public String toFullString() {
return "DefaultServiceInstance{" +
", serviceName='" + serviceName + '\'' +
", host='" + host + '\'' +
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistryDirectory.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistryDirectory.java
index 505142e..3b9e05e 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistryDirectory.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistryDirectory.java
@@ -103,7 +103,7 @@ public class ServiceDiscoveryRegistryDirectory<T> extends
DynamicDirectory<T> im
}
/**
- * This implementation wants to make sure all application names related to
serviceListener received address notification.
+ * This implementation makes sure all application names related to
serviceListener received address notification.
*
* FIXME, make sure deprecated "interface-application" mapping item be
cleared in time.
*/
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/metadata/store/InMemoryWritableMetadataService.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/metadata/store/InMemoryWritableMetadataService.java
index f7354ab..6b27202 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/metadata/store/InMemoryWritableMetadataService.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/metadata/store/InMemoryWritableMetadataService.java
@@ -317,6 +317,11 @@ public class InMemoryWritableMetadataService implements
WritableMetadataService
}
@Override
+ public Map<String, Set<String>> getCachedMapping() {
+ return serviceToAppsMapping;
+ }
+
+ @Override
public void setMetadataServiceURL(URL url) {
this.metadataServiceURL = url;
}
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/DefaultMigrationAddressComparator.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/DefaultMigrationAddressComparator.java
index a9c7b8c..9700214 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/DefaultMigrationAddressComparator.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/DefaultMigrationAddressComparator.java
@@ -27,31 +27,44 @@ import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.cluster.ClusterInvoker;
import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
public class DefaultMigrationAddressComparator implements
MigrationAddressComparator {
private static final Logger logger =
LoggerFactory.getLogger(DefaultMigrationAddressComparator.class);
private static final String MIGRATION_THRESHOLD =
"dubbo.application.migration.threshold";
- private static final String DEFAULT_THRESHOLD_STRING = "0.8";
- private static final float DEFAULT_THREAD = 0.8f;
+ private static final String DEFAULT_THRESHOLD_STRING = "1.0";
+ private static final float DEFAULT_THREAD = 1.0f;
+
+ public static final String OLD_ADDRESS_SIZE = "OLD_ADDRESS_SIZE";
+ public static final String NEW_ADDRESS_SIZE = "NEW_ADDRESS_SIZE";
private static final WritableMetadataService localMetadataService =
WritableMetadataService.getDefaultExtension();
+ private Map<String, Map<String, Integer>> serviceMigrationData = new
ConcurrentHashMap<>();
+
@Override
public <T> boolean shouldMigrate(ClusterInvoker<T>
serviceDiscoveryInvoker, ClusterInvoker<T> invoker, MigrationRule rule) {
+ Map<String, Integer> migrationData =
serviceMigrationData.computeIfAbsent(invoker.getUrl().getDisplayServiceKey(),
_k -> new ConcurrentHashMap<>());
+
if (!serviceDiscoveryInvoker.hasProxyInvokers()) {
+ migrationData.put(OLD_ADDRESS_SIZE, getAddressSize(invoker));
+ migrationData.put(NEW_ADDRESS_SIZE, -1);
logger.info("No instance address available, stop compare.");
return false;
}
if (!invoker.hasProxyInvokers()) {
+ migrationData.put(OLD_ADDRESS_SIZE, -1);
+ migrationData.put(NEW_ADDRESS_SIZE,
getAddressSize(serviceDiscoveryInvoker));
logger.info("No interface address available, stop compare.");
return true;
}
- List<Invoker<T>> invokers1 =
serviceDiscoveryInvoker.getDirectory().getAllInvokers();
- List<Invoker<T>> invokers2 = invoker.getDirectory().getAllInvokers();
+ int newAddressSize = getAddressSize(serviceDiscoveryInvoker);
+ int oldAddressSize = getAddressSize(invoker);
- int newAddressSize = CollectionUtils.isNotEmpty(invokers1) ?
invokers1.size() : 0;
- int oldAddressSize = CollectionUtils.isNotEmpty(invokers2) ?
invokers2.size() : 0;
+ migrationData.put(OLD_ADDRESS_SIZE, oldAddressSize);
+ migrationData.put(NEW_ADDRESS_SIZE, newAddressSize);
String rawThreshold = null;
String serviceKey = invoker.getUrl().getDisplayServiceKey();
@@ -82,4 +95,17 @@ public class DefaultMigrationAddressComparator implements
MigrationAddressCompar
}
return false;
}
+
+ private <T> int getAddressSize(ClusterInvoker<T> invoker) {
+ if (invoker == null) {
+ return -1;
+ }
+ List<Invoker<T>> invokers = invoker.getDirectory().getAllInvokers();
+ return CollectionUtils.isNotEmpty(invokers) ? invokers.size() : 0;
+ }
+
+ public Map<String, Integer> getAddressSize(String displayServiceKey) {
+ return serviceMigrationData.get(displayServiceKey);
+ }
+
}
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationAddressComparator.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationAddressComparator.java
index 57906ca..32fedab 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationAddressComparator.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationAddressComparator.java
@@ -20,7 +20,10 @@ import org.apache.dubbo.common.extension.SPI;
import org.apache.dubbo.registry.client.migration.model.MigrationRule;
import org.apache.dubbo.rpc.cluster.ClusterInvoker;
+import java.util.Map;
+
@SPI
public interface MigrationAddressComparator {
<T> boolean shouldMigrate(ClusterInvoker<T> serviceDiscoveryInvoker,
ClusterInvoker<T> invoker, MigrationRule rule);
+ Map<String, Integer> getAddressSize(String displayServiceKey);
}
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationInvoker.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationInvoker.java
index 85e216a..bc23bf7 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationInvoker.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/migration/MigrationInvoker.java
@@ -20,6 +20,7 @@ import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.logger.Logger;
import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.status.reporter.FrameworkStatusReporter;
import org.apache.dubbo.common.utils.CollectionUtils;
import org.apache.dubbo.common.utils.StringUtils;
import org.apache.dubbo.registry.Registry;
@@ -40,6 +41,7 @@ import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
+import static
org.apache.dubbo.common.status.reporter.FrameworkStatusReporter.createConsumptionReport;
import static org.apache.dubbo.rpc.cluster.Constants.REFER_KEY;
public class MigrationInvoker<T> implements MigrationClusterInvoker<T> {
@@ -59,7 +61,6 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
private volatile MigrationRule rule;
private volatile boolean migrated;
-
private static final ScheduledExecutorService scheduler =
Executors.newSingleThreadScheduledExecutor();
public MigrationInvoker(RegistryProtocol registryProtocol,
@@ -87,6 +88,11 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
this.type = type;
this.url = url;
this.consumerUrl = consumerUrl;
+
+ ConsumerModel consumerModel =
ApplicationModel.getConsumerModel(consumerUrl.getServiceKey());
+ if (consumerModel != null) {
+
consumerModel.getServiceMetadata().addAttribute("currentClusterInvoker", this);
+ }
}
public ClusterInvoker<T> getInvoker() {
@@ -139,9 +145,16 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
@Override
public void fallbackToInterfaceInvoker() {
+ migrated = false;
refreshInterfaceInvoker();
setListener(invoker, () -> {
- this.destroyServiceDiscoveryInvoker();
+ if (!migrated) {
+ migrated = true;
+ this.destroyServiceDiscoveryInvoker();
+ FrameworkStatusReporter.reportConsumptionStatus(
+
createConsumptionReport(consumerUrl.getServiceInterface(),
consumerUrl.getVersion(), consumerUrl.getGroup(), "interface")
+ );
+ }
});
}
@@ -151,6 +164,8 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
fallbackToInterfaceInvoker();
return;
}
+
+ migrated = false;
if (!forceMigrate) {
refreshServiceDiscoveryInvoker();
refreshInterfaceInvoker();
@@ -163,7 +178,13 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
} else {
refreshServiceDiscoveryInvoker();
setListener(serviceDiscoveryInvoker, () -> {
- this.destroyInterfaceInvoker();
+ if (!migrated) {
+ migrated = true;
+ this.destroyInterfaceInvoker();
+ FrameworkStatusReporter.reportConsumptionStatus(
+
createConsumptionReport(consumerUrl.getServiceInterface(),
consumerUrl.getVersion(), consumerUrl.getGroup(), "app")
+ );
+ }
});
}
}
@@ -228,6 +249,10 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
if (serviceDiscoveryInvoker != null) {
serviceDiscoveryInvoker.destroy();
}
+ ConsumerModel consumerModel =
ApplicationModel.getConsumerModel(consumerUrl.getServiceKey());
+ if (consumerModel != null) {
+
consumerModel.getServiceMetadata().getAttributeMap().remove("currentClusterInvoker");
+ }
}
@Override
@@ -329,6 +354,9 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
if (invoker.getDirectory().isNotificationReceived()) {
destroyInterfaceInvoker();
migrated = true;
+ FrameworkStatusReporter.reportConsumptionStatus(
+
createConsumptionReport(consumerUrl.getServiceInterface(),
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_app")
+ );
}
});
} else {
@@ -337,6 +365,9 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
if
(serviceDiscoveryInvoker.getDirectory().isNotificationReceived()) {
destroyServiceDiscoveryInvoker();
migrated = true;
+ FrameworkStatusReporter.reportConsumptionStatus(
+
createConsumptionReport(consumerUrl.getServiceInterface(),
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_interface")
+ );
}
});
}
@@ -353,9 +384,6 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
serviceDiscoveryInvoker.destroy();
serviceDiscoveryInvoker = null;
}
-
- updateConsumerModel(currentAvailableInvoker, serviceDiscoveryInvoker);
- migrated = true;
}
// protected synchronized void
discardServiceDiscoveryInvokerAddress(ClusterInvoker<T>
serviceDiscoveryInvoker) {
@@ -405,9 +433,6 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
invoker.destroy();
invoker = null;
}
-
- updateConsumerModel(currentAvailableInvoker, invoker);
- migrated = true;
}
//
@@ -447,16 +472,4 @@ public class MigrationInvoker<T> implements
MigrationClusterInvoker<T> {
public boolean checkInvokerAvailable(ClusterInvoker<T> invoker) {
return invoker != null && !invoker.isDestroyed() &&
invoker.isAvailable();
}
-
- private void updateConsumerModel(ClusterInvoker<?> workingInvoker,
ClusterInvoker<?> backInvoker) {
- ConsumerModel consumerModel =
ApplicationModel.getConsumerModel(consumerUrl.getServiceKey());
- if (consumerModel != null) {
- if (workingInvoker != null) {
-
consumerModel.getServiceMetadata().addAttribute("currentClusterInvoker",
workingInvoker);
- }
- if (backInvoker != null) {
-
consumerModel.getServiceMetadata().addAttribute("backupClusterInvoker",
backInvoker);
- }
- }
- }
}
diff --git
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/RegistryDirectory.java
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/RegistryDirectory.java
index efce5d0..00fb54a 100644
---
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/RegistryDirectory.java
+++
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/RegistryDirectory.java
@@ -236,10 +236,10 @@ public class RegistryDirectory<T> extends
DynamicDirectory<T> implements NotifyL
} catch (Exception e) {
logger.warn("destroyUnusedInvokers error. ", e);
}
- }
- // notify invokers refreshed
- this.invokersChanged();
+ // notify invokers refreshed
+ this.invokersChanged();
+ }
}
private List<Invoker<T>> toMergeInvokerList(List<Invoker<T>> invokers) {