diff --git a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistry.java b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistry.java index 1a79ea2325..65f50af4d8 100644 --- a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistry.java +++ b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/client/ServiceDiscoveryRegistry.java @@ -323,20 +323,24 @@ public class ServiceDiscoveryRegistry implements Registry { serviceToAppsMapping.put(protocolServiceKey, serviceNamesKey); // register ServiceInstancesChangedListener - ServiceInstancesChangedListener serviceListener = serviceListeners.computeIfAbsent(serviceNamesKey, - k -> new ServiceInstancesChangedListener(serviceNames, serviceDiscovery)); - serviceListener.setUrl(url); - listener.addServiceListener(serviceListener); - - serviceNames.forEach(serviceName -> { - List serviceInstances = serviceDiscovery.getInstances(serviceName); - serviceListener.onEvent(new ServiceInstancesChangedEvent(serviceName, serviceInstances)); + ServiceInstancesChangedListener serviceListener = serviceListeners.computeIfAbsent(serviceNamesKey, k -> { + ServiceInstancesChangedListener serviceInstancesChangedListener = new ServiceInstancesChangedListener(serviceNames, serviceDiscovery); + serviceInstancesChangedListener.setUrl(url); + serviceNames.forEach(serviceName -> { + List serviceInstances = serviceDiscovery.getInstances(serviceName); + if (CollectionUtils.isNotEmpty(serviceInstances)) { + serviceInstancesChangedListener.onEvent(new ServiceInstancesChangedEvent(serviceName, serviceInstances)); + } + }); + return serviceInstancesChangedListener; }); - listener.notify(serviceListener.getUrls(protocolServiceKey)); - + serviceListener.setUrl(url); + listener.addServiceListener(serviceListener); serviceListener.addListener(protocolServiceKey, listener); registerServiceInstancesChangedListener(url, serviceListener); + + listener.notify(serviceListener.getUrls(protocolServiceKey)); } /** 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 9af31cd6b2..20dc035b22 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 @@ -26,12 +26,14 @@ import org.apache.dubbo.registry.client.migration.model.MigrationStep; import org.apache.dubbo.registry.integration.DynamicDirectory; import org.apache.dubbo.registry.integration.RegistryProtocol; import org.apache.dubbo.rpc.Invocation; +import org.apache.dubbo.rpc.Invoker; import org.apache.dubbo.rpc.Result; import org.apache.dubbo.rpc.RpcException; import org.apache.dubbo.rpc.cluster.Cluster; import org.apache.dubbo.rpc.cluster.ClusterInvoker; import org.apache.dubbo.rpc.cluster.Directory; +import java.util.List; import java.util.Set; import static org.apache.dubbo.rpc.cluster.Constants.REFER_KEY; @@ -279,13 +281,14 @@ public class MigrationInvoker implements MigrationClusterInvoker { } } - protected synchronized void discardServiceDiscoveryInvokerAddress(ClusterInvoker serviceDiscoveryInvoker) { + protected synchronized void discardServiceDiscoveryInvokerAddress(ClusterInvoker serviceDiscoveryInvoker) { if (this.invoker != null) { this.currentAvailableInvoker = this.invoker; } if (serviceDiscoveryInvoker != null) { if (logger.isDebugEnabled()) { - logger.debug("Discarding instance addresses, total size " + serviceDiscoveryInvoker.getDirectory().getAllInvokers().size()); + List> invokers = serviceDiscoveryInvoker.getDirectory().getAllInvokers(); + logger.debug("Discarding instance addresses, total size " + (invokers == null ? 0 : invokers.size())); } // serviceDiscoveryInvoker.getDirectory().discordAddresses(); } @@ -298,6 +301,8 @@ public class MigrationInvoker implements MigrationClusterInvoker { logger.debug("Re-subscribing instance addresses, current interface " + type.getName()); } serviceDiscoveryInvoker = registryProtocol.getServiceDiscoveryInvoker(cluster, registry, type, url); + } else { + ((DynamicDirectory)serviceDiscoveryInvoker.getDirectory()).markInvokersChanged(); } } @@ -310,6 +315,8 @@ public class MigrationInvoker implements MigrationClusterInvoker { } invoker = registryProtocol.getInvoker(cluster, registry, type, url); + } else { + ((DynamicDirectory)invoker.getDirectory()).markInvokersChanged(); } } @@ -331,7 +338,8 @@ public class MigrationInvoker implements MigrationClusterInvoker { } if (invoker != null) { if (logger.isDebugEnabled()) { - logger.debug("Discarding interface addresses, total address size " + invoker.getDirectory().getAllInvokers().size()); + List> invokers = invoker.getDirectory().getAllInvokers(); + logger.debug("Discarding interface addresses, total address size " + (invokers == null ? 0 : invokers.size())); } //invoker.getDirectory().discordAddresses(); } diff --git a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/DynamicDirectory.java b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/DynamicDirectory.java index 8d2ed8b122..56dc3a985f 100644 --- a/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/DynamicDirectory.java +++ b/dubbo-registry/dubbo-registry-api/src/main/java/org/apache/dubbo/registry/integration/DynamicDirectory.java @@ -260,17 +260,27 @@ public abstract class DynamicDirectory extends AbstractDirectory implement } private volatile InvokersChangedListener invokersChangedListener; + private volatile boolean invokersChanged; - public void setInvokersChangedListener(InvokersChangedListener listener) { + public synchronized void setInvokersChangedListener(InvokersChangedListener listener) { this.invokersChangedListener = listener; - invokersChanged(); + if (invokersChangedListener != null && invokersChanged) { + invokersChangedListener.onChange(); + invokersChanged = false; + } } - protected void invokersChanged() { - if (invokersChangedListener != null) { + protected synchronized void invokersChanged() { + invokersChanged = true; + if (invokersChangedListener != null && invokersChanged) { invokersChangedListener.onChange(); + invokersChanged = false; } } + public synchronized void markInvokersChanged() { + this.invokersChanged = true; + } + protected abstract void destroyAllInvokers(); }