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 4bfeeb2  Continue to notify address on metadata failure (#7796)
4bfeeb2 is described below

commit 4bfeeb2b47284ae04f1fe340ac4c29c096dfb260
Author: ken.lj <[email protected]>
AuthorDate: Wed May 19 10:43:11 2021 +0800

    Continue to notify address on metadata failure (#7796)
---
 .../src/main/resources/dubbo-migration.yaml        |   2 +-
 .../src/main/resources/spring/dubbo-consumer.xml   |   2 +-
 .../listener/ServiceInstancesChangedListener.java  |   8 +-
 .../client/migration/MigrationInvoker.java         | 145 ++++++++++-----------
 4 files changed, 78 insertions(+), 79 deletions(-)

diff --git 
a/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/dubbo-migration.yaml
 
b/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/dubbo-migration.yaml
index b33d32c..3f55a0b 100644
--- 
a/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/dubbo-migration.yaml
+++ 
b/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/dubbo-migration.yaml
@@ -18,5 +18,5 @@
 #
 
 key: demo-consumer
-step: FORCE_APPLICATION
+step: APPLICATION_FIRST
 threshold: 0.1
\ No newline at end of file
diff --git 
a/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/spring/dubbo-consumer.xml
 
b/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/spring/dubbo-consumer.xml
index 539a645..481c9b4 100644
--- 
a/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/spring/dubbo-consumer.xml
+++ 
b/dubbo-demo/dubbo-demo-xml/dubbo-demo-xml-consumer/src/main/resources/spring/dubbo-consumer.xml
@@ -27,7 +27,7 @@
 
     <!--    <dubbo:metadata-report address="zookeeper://127.0.0.1:2181"/>-->
 
-    <dubbo:registry 
address="zookeeper://127.0.0.1:2181?registry-type=service"/>
+    <dubbo:registry address="zookeeper://127.0.0.1:2181"/>
 
     <dubbo:reference id="demoService" check="false"
                      interface="org.apache.dubbo.demo.DemoService"/>
diff --git 
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/event/listener/ServiceInstancesChangedListener.java
 
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/event/listener/ServiceInstancesChangedListener.java
index 1517ed7..b8bac8e 100644
--- 
a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/event/listener/ServiceInstancesChangedListener.java
+++ 
b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/event/listener/ServiceInstancesChangedListener.java
@@ -149,9 +149,9 @@ public class ServiceInstancesChangedListener implements 
ConditionalEventListener
                 retryFuture = scheduler.schedule(new 
AddressRefreshRetryTask(retryPermission), 10000, TimeUnit.MILLISECONDS);
                 logger.warn("Address refresh try task submitted.");
             }
-            logger.warn("Address refresh failed because of Metadata Server 
failure, wait for retry or new address refresh event.");
-            this.revisionToMetadata = newRevisionToMetadata;
-            return;
+//            logger.warn("Address refresh failed because of Metadata Server 
failure, wait for retry or new address refresh event.");
+//            this.revisionToMetadata = newRevisionToMetadata;
+//            return;
         }
 
         this.revisionToMetadata = newRevisionToMetadata;
@@ -244,7 +244,7 @@ public class ServiceInstancesChangedListener implements 
ConditionalEventListener
         if (event instanceof RetryServiceInstancesChangedEvent) {
             RetryServiceInstancesChangedEvent retryEvent = 
(RetryServiceInstancesChangedEvent) event;
             logger.warn("Received address refresh retry event, " + 
retryEvent.getFailureRecordTime());
-            if (retryEvent.getFailureRecordTime() < lastRefreshTime) {
+            if (retryEvent.getFailureRecordTime() < lastRefreshTime && 
!hasEmptyMetadata(revisionToMetadata)) {
                 logger.warn("Ignore retry event, event time: " + 
retryEvent.getFailureRecordTime() + ", last refresh time: " + lastRefreshTime);
                 return true;
             }
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 bc23bf7..f1f5ef3 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
@@ -17,6 +17,7 @@
 package org.apache.dubbo.registry.client.migration;
 
 import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.config.ConfigurationUtils;
 import org.apache.dubbo.common.extension.ExtensionLoader;
 import org.apache.dubbo.common.logger.Logger;
 import org.apache.dubbo.common.logger.LoggerFactory;
@@ -40,12 +41,14 @@ import org.apache.dubbo.rpc.model.ConsumerModel;
 import java.util.Set;
 import java.util.concurrent.Executors;
 import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 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> {
     private Logger logger = LoggerFactory.getLogger(MigrationInvoker.class);
+    private static final String MIGRATION_DELAY_KEY = 
"dubbo.application.migration.delay";
 
     private URL url;
     private URL consumerUrl;
@@ -169,12 +172,17 @@ public class MigrationInvoker<T> implements 
MigrationClusterInvoker<T> {
         if (!forceMigrate) {
             refreshServiceDiscoveryInvoker();
             refreshInterfaceInvoker();
-            setListener(invoker, () -> {
-                this.compareAddresses(serviceDiscoveryInvoker, invoker);
-            });
-            setListener(serviceDiscoveryInvoker, () -> {
-                this.compareAddresses(serviceDiscoveryInvoker, invoker);
-            });
+            // By the time the task gets scheduled, the address notifications 
are expected to be finished for both address
+            // types. Otherwise, migration task will just pick interface 
invoker.
+            if (!migrated) {
+                scheduler.schedule(new MigrationTask(), getDelay(), 
TimeUnit.MILLISECONDS);
+                setListener(invoker, () -> {
+                    this.setAvailableInvoker(serviceDiscoveryInvoker, invoker);
+                });
+                setListener(serviceDiscoveryInvoker, () -> {
+                    this.setAvailableInvoker(serviceDiscoveryInvoker, invoker);
+                });
+            }
         } else {
             refreshServiceDiscoveryInvoker();
             setListener(serviceDiscoveryInvoker, () -> {
@@ -189,6 +197,51 @@ public class MigrationInvoker<T> implements 
MigrationClusterInvoker<T> {
         }
     }
 
+    private int getDelay() {
+        int delay = 60000;
+        String delayStr = ConfigurationUtils.getProperty(MIGRATION_DELAY_KEY);
+        if (StringUtils.isEmpty(delayStr)) {
+            return delay;
+        }
+
+        try {
+           delay = Integer.parseInt(delayStr);
+        } catch (Exception e) {
+            logger.warn("Invalid migration delay param " + delayStr);
+        }
+        return delay;
+    }
+
+    private class MigrationTask implements Runnable {
+        @Override
+        public void run() {
+            if (migrated) {
+                return;
+            }
+
+            Set<MigrationAddressComparator> detectors = 
ExtensionLoader.getExtensionLoader(MigrationAddressComparator.class).getSupportedExtensionInstances();
+            if (CollectionUtils.isEmpty(detectors)) {
+                migrated = true;
+                destroyInterfaceInvoker();
+                return;
+            }
+
+            if (detectors.stream().allMatch(comparator -> 
comparator.shouldMigrate(serviceDiscoveryInvoker, invoker, rule))) {
+                destroyInterfaceInvoker();
+                FrameworkStatusReporter.reportConsumptionStatus(
+                        
createConsumptionReport(consumerUrl.getServiceInterface(), 
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_app")
+                );
+            } else {
+                // by default
+                destroyServiceDiscoveryInvoker();
+                FrameworkStatusReporter.reportConsumptionStatus(
+                        
createConsumptionReport(consumerUrl.getServiceInterface(), 
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_interface")
+                );
+            }
+            migrated = true;
+        }
+    }
+
     private boolean checkMigratingConditionMatch(URL consumerUrl) {
         Set<PreMigratingConditionChecker> checkers = 
ExtensionLoader.getExtensionLoader(PreMigratingConditionChecker.class).getSupportedExtensionInstances();
         if (CollectionUtils.isNotEmpty(checkers)) {
@@ -332,44 +385,22 @@ public class MigrationInvoker<T> implements 
MigrationClusterInvoker<T> {
     private volatile boolean invokersChanged;
 
     /**
-     * Need to know which invoker change triggered this compare.
+     * Set the available invoker before migration is determined. Interface 
invoker goes first.
+     * @param serviceDiscoveryInvoker
+     * @param invoker
      */
-    private synchronized void compareAddresses(ClusterInvoker<T> 
serviceDiscoveryInvoker, ClusterInvoker<T> invoker) {
+    private void setAvailableInvoker(ClusterInvoker<T> 
serviceDiscoveryInvoker, ClusterInvoker<T> invoker) {
         this.invokersChanged = true;
-        if (migrated) {
-            return;
-        }
-        if (serviceDiscoveryInvoker == null || 
serviceDiscoveryInvoker.isDestroyed()) {
-            currentAvailableInvoker = invoker;
-            return;
-        } else if (invoker == null || invoker.isDestroyed()) {
-            currentAvailableInvoker = serviceDiscoveryInvoker;
-            return;
-        }
 
-        Set<MigrationAddressComparator> detectors = 
ExtensionLoader.getExtensionLoader(MigrationAddressComparator.class).getSupportedExtensionInstances();
-        if (detectors != null && detectors.stream().allMatch(migrationDetector 
-> migrationDetector.shouldMigrate(serviceDiscoveryInvoker, invoker, rule))) {
-            logger.info("serviceKey:" + invoker.getUrl().getServiceKey() + " 
switch to APP Level address");
-            scheduler.submit(() -> {
-                if (invoker.getDirectory().isNotificationReceived()) {
-                    destroyInterfaceInvoker();
-                    migrated = true;
-                    FrameworkStatusReporter.reportConsumptionStatus(
-                            
createConsumptionReport(consumerUrl.getServiceInterface(), 
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_app")
-                    );
-                }
-            });
-        } else {
-            logger.info("serviceKey:" + 
serviceDiscoveryInvoker.getUrl().getServiceKey() + " switch to Service Level 
address");
-            scheduler.submit(() -> {
-                if 
(serviceDiscoveryInvoker.getDirectory().isNotificationReceived()) {
-                    destroyServiceDiscoveryInvoker();
-                    migrated = true;
-                    FrameworkStatusReporter.reportConsumptionStatus(
-                            
createConsumptionReport(consumerUrl.getServiceInterface(), 
consumerUrl.getVersion(), consumerUrl.getGroup(), "app_interface")
-                    );
-                }
-            });
+        if (currentAvailableInvoker == null) {
+            if (invoker != null && !invoker.isDestroyed() && 
invoker.hasProxyInvokers()) {
+                currentAvailableInvoker = invoker;
+            } else if (serviceDiscoveryInvoker != null && 
!serviceDiscoveryInvoker.isDestroyed() && 
serviceDiscoveryInvoker.hasProxyInvokers()) {
+                currentAvailableInvoker = serviceDiscoveryInvoker;
+            }
+            if (currentAvailableInvoker == null) {
+                currentAvailableInvoker = invoker;
+            }
         }
     }
 
@@ -386,20 +417,6 @@ public class MigrationInvoker<T> implements 
MigrationClusterInvoker<T> {
         }
     }
 
-//    protected synchronized void 
discardServiceDiscoveryInvokerAddress(ClusterInvoker<T> 
serviceDiscoveryInvoker) {
-//        if (this.invoker != null) {
-//            this.currentAvailableInvoker = this.invoker;
-//            updateConsumerModel(currentAvailableInvoker, 
serviceDiscoveryInvoker);
-//        }
-//        if (serviceDiscoveryInvoker != null && 
!serviceDiscoveryInvoker.isDestroyed()) {
-//            if (logger.isDebugEnabled()) {
-//                List<Invoker<T>> invokers = 
serviceDiscoveryInvoker.getDirectory().getAllInvokers();
-//                logger.debug("Discarding instance addresses, total size " + 
(invokers == null ? 0 : invokers.size()));
-//            }
-//            serviceDiscoveryInvoker.getDirectory().discordAddresses();
-//        }
-//    }
-
     protected void refreshServiceDiscoveryInvoker() {
         clearListener(serviceDiscoveryInvoker);
         if (needRefresh(serviceDiscoveryInvoker)) {
@@ -435,24 +452,6 @@ public class MigrationInvoker<T> implements 
MigrationClusterInvoker<T> {
         }
     }
 
-//
-//    protected synchronized void 
discardInterfaceInvokerAddress(ClusterInvoker<T> invoker) {
-//        if (this.serviceDiscoveryInvoker != null) {
-//            this.currentAvailableInvoker = this.serviceDiscoveryInvoker;
-//            updateConsumerModel(currentAvailableInvoker, invoker);
-//        }
-//        if (invoker != null && !invoker.isDestroyed()) {
-//            if (logger.isDebugEnabled()) {
-//                List<Invoker<T>> invokers = 
invoker.getDirectory().getAllInvokers();
-//                logger.debug("Discarding interface addresses, total address 
size " + (invokers == null ? 0 : invokers.size()));
-//            }
-//            invoker.getDirectory().discordAddresses();
-////            if (invokerDestroyStatus == null) {
-////                invokerDestroyStatus = executorService.schedule(new 
InvokerDestroyTask(), destroyInterval, TimeUnit.MILLISECONDS);
-////            }
-//        }
-//    }
-
     private void clearListener(ClusterInvoker<T> invoker) {
         if (invoker == null) return;
         DynamicDirectory<T> directory = (DynamicDirectory<T>) 
invoker.getDirectory();

Reply via email to