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 1abd3831f9..7740ec87dd 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 @@ -22,6 +22,7 @@ import org.apache.dubbo.common.logger.ErrorTypeAwareLogger; import org.apache.dubbo.common.logger.LoggerFactory; import org.apache.dubbo.common.utils.CollectionUtils; import org.apache.dubbo.common.utils.ConcurrentHashMapUtils; +import org.apache.dubbo.common.utils.ConcurrentHashSet; import org.apache.dubbo.metadata.AbstractServiceNameMapping; import org.apache.dubbo.metadata.MappingChangedEvent; import org.apache.dubbo.metadata.MappingListener; @@ -37,11 +38,13 @@ import org.apache.dubbo.rpc.model.ApplicationModel; import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; +import java.util.stream.Collectors; import static org.apache.dubbo.common.constants.CommonConstants.INTERFACE_KEY; import static org.apache.dubbo.common.constants.CommonConstants.PROVIDER_SIDE; @@ -79,7 +82,7 @@ public class ServiceDiscoveryRegistry extends FailbackRegistry { /* apps - listener */ private final Map serviceListeners = new ConcurrentHashMap<>(); - private final Map mappingListeners = new ConcurrentHashMap<>(); + private final Map> mappingListeners = new ConcurrentHashMap<>(); /* This lock has the same scope and lifecycle as its corresponding instance listener. It's used to make sure that only one interface mapping to the same app list can do subscribe or unsubscribe at the same moment. And the lock should be destroyed when listener destroying its corresponding instance listener. @@ -112,7 +115,7 @@ public class ServiceDiscoveryRegistry extends FailbackRegistry { */ protected ServiceDiscovery createServiceDiscovery(URL registryURL) { return getServiceDiscovery(registryURL.addParameter(INTERFACE_KEY, ServiceDiscovery.class.getName()) - .removeParameter(REGISTRY_TYPE_KEY)); + .removeParameter(REGISTRY_TYPE_KEY)); } /** @@ -206,7 +209,10 @@ public class ServiceDiscoveryRegistry extends FailbackRegistry { try { MappingListener mappingListener = new DefaultMappingListener(url, mappingByUrl, listener); mappingByUrl = serviceNameMapping.getAndListen(this.getUrl(), url, mappingListener); - mappingListeners.put(url.getProtocolServiceKey(), mappingListener); + synchronized (mappingListeners) { + mappingListeners.computeIfAbsent(url.getProtocolServiceKey(), (k) -> new ConcurrentHashSet<>()) + .add(mappingListener); + } } catch (Exception e) { logger.warn(INTERNAL_ERROR, "", "", "Cannot find app mapping for service " + url.getServiceInterface() + ", will not migrate.", e); } @@ -248,8 +254,23 @@ public class ServiceDiscoveryRegistry extends FailbackRegistry { serviceDiscovery.unsubscribe(url, listener); String protocolServiceKey = url.getProtocolServiceKey(); Set serviceNames = serviceNameMapping.getMapping(url); - if (mappingListeners.get(protocolServiceKey) != null) { - serviceNameMapping.stopListen(url, mappingListeners.remove(protocolServiceKey)); + + synchronized (mappingListeners) { + Set keyedListeners = mappingListeners.get(protocolServiceKey); + if (keyedListeners != null) { + List matched = keyedListeners.stream() + .filter(mappingListener -> + mappingListener instanceof DefaultMappingListener + && (Objects.equals(((DefaultMappingListener) mappingListener).getListener(), listener))) + .collect(Collectors.toList()); + for (MappingListener mappingListener : matched) { + serviceNameMapping.stopListen(url, mappingListener); + keyedListeners.remove(mappingListener); + } + if (keyedListeners.isEmpty()) { + mappingListeners.remove(protocolServiceKey, Collections.emptySet()); + } + } } if (CollectionUtils.isNotEmpty(serviceNames)) { String serviceNamesKey = toStringKeys(serviceNames); @@ -430,6 +451,10 @@ public class ServiceDiscoveryRegistry extends FailbackRegistry { } } + protected NotifyListener getListener() { + return listener; + } + @Override public void stop() { stopped = true;