diff --git a/dubbo-common/pom.xml b/dubbo-common/pom.xml
index b424661be6..f32317c093 100644
--- a/dubbo-common/pom.xml
+++ b/dubbo-common/pom.xml
@@ -63,6 +63,10 @@
com.alibaba
fastjson
+
+ com.google.code.gson
+ gson
+
commons-io
commons-io
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 0000000000..6d32e9ffe6
--- /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 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 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 6919652d5b..2f68774ca2 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 bac064ad83..f524fcf7d3 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 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 a233910473..288712ee9f 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 1782ef7788..222da5fe97 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 removeCachedMapping(String serviceKey);
+ Map> 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 1088b7b689..90d8127f62 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 discoveries = new ConcurrentHashMap<>();
+ private static final ConcurrentMap discoveries = new ConcurrentHashMap<>();
+
+ public static List 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 1bcf284d46..10fe957487 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 505142e7eb..3b9e05e943 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 extends DynamicDirectory 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 f7354abb78..6b27202dbb 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
@@ -316,6 +316,11 @@ public class InMemoryWritableMetadataService implements WritableMetadataService
return serviceToAppsMapping.remove(serviceKey);
}
+ @Override
+ public Map> 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 a9c7b8ce31..97002146d5 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> serviceMigrationData = new ConcurrentHashMap<>();
+
@Override
public boolean shouldMigrate(ClusterInvoker serviceDiscoveryInvoker, ClusterInvoker invoker, MigrationRule rule) {
+ Map 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> invokers1 = serviceDiscoveryInvoker.getDirectory().getAllInvokers();
- List> 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 int getAddressSize(ClusterInvoker invoker) {
+ if (invoker == null) {
+ return -1;
+ }
+ List> invokers = invoker.getDirectory().getAllInvokers();
+ return CollectionUtils.isNotEmpty(invokers) ? invokers.size() : 0;
+ }
+
+ public Map 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 57906ca5c0..32fedab926 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 {
boolean shouldMigrate(ClusterInvoker serviceDiscoveryInvoker, ClusterInvoker invoker, MigrationRule rule);
+ Map 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 85e216a229..bc23bf7a36 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 implements MigrationClusterInvoker {
@@ -59,7 +61,6 @@ public class MigrationInvoker implements MigrationClusterInvoker {
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 implements MigrationClusterInvoker {
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 getInvoker() {
@@ -139,9 +145,16 @@ public class MigrationInvoker implements MigrationClusterInvoker {
@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 implements MigrationClusterInvoker {
fallbackToInterfaceInvoker();
return;
}
+
+ migrated = false;
if (!forceMigrate) {
refreshServiceDiscoveryInvoker();
refreshInterfaceInvoker();
@@ -163,7 +178,13 @@ public class MigrationInvoker implements MigrationClusterInvoker {
} 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 implements MigrationClusterInvoker {
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 implements MigrationClusterInvoker {
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 implements MigrationClusterInvoker {
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 implements MigrationClusterInvoker {
serviceDiscoveryInvoker.destroy();
serviceDiscoveryInvoker = null;
}
-
- updateConsumerModel(currentAvailableInvoker, serviceDiscoveryInvoker);
- migrated = true;
}
// protected synchronized void discardServiceDiscoveryInvokerAddress(ClusterInvoker serviceDiscoveryInvoker) {
@@ -405,9 +433,6 @@ public class MigrationInvoker implements MigrationClusterInvoker {
invoker.destroy();
invoker = null;
}
-
- updateConsumerModel(currentAvailableInvoker, invoker);
- migrated = true;
}
//
@@ -447,16 +472,4 @@ public class MigrationInvoker implements MigrationClusterInvoker {
public boolean checkInvokerAvailable(ClusterInvoker 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 efce5d0b59..00fb54a309 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 extends DynamicDirectory implements NotifyL
} catch (Exception e) {
logger.warn("destroyUnusedInvokers error. ", e);
}
- }
- // notify invokers refreshed
- this.invokersChanged();
+ // notify invokers refreshed
+ this.invokersChanged();
+ }
}
private List> toMergeInvokerList(List> invokers) {