[3.0] Fix Nacos aggregate listen (#10468)
* [3.0] Fix Nacos aggregate listen * fix empty * rename
This commit is contained in:
parent
51529c34c5
commit
43c4033e11
|
|
@ -0,0 +1,70 @@
|
|||
/*
|
||||
* 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.registry.nacos;
|
||||
|
||||
import org.apache.dubbo.common.utils.ConcurrentHashSet;
|
||||
import org.apache.dubbo.registry.NotifyListener;
|
||||
|
||||
import com.alibaba.nacos.api.naming.pojo.Instance;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
public class NacosAggregateListener {
|
||||
private final NotifyListener notifyListener;
|
||||
private final Set<String> serviceNames = new ConcurrentHashSet<>();
|
||||
private final Map<String, List<Instance>> serviceInstances = new ConcurrentHashMap<>();
|
||||
|
||||
public NacosAggregateListener(NotifyListener notifyListener) {
|
||||
this.notifyListener = notifyListener;
|
||||
}
|
||||
|
||||
public List<Instance> saveAndAggregateAllInstances(String serviceName, List<Instance> instances) {
|
||||
serviceNames.add(serviceName);
|
||||
serviceInstances.put(serviceName, instances);
|
||||
return serviceInstances.values().stream().flatMap(List::stream).collect(Collectors.toList());
|
||||
}
|
||||
|
||||
public NotifyListener getNotifyListener() {
|
||||
return notifyListener;
|
||||
}
|
||||
|
||||
public Set<String> getServiceNames() {
|
||||
return serviceNames;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean equals(Object o) {
|
||||
if (this == o) {
|
||||
return true;
|
||||
}
|
||||
if (o == null || getClass() != o.getClass()) {
|
||||
return false;
|
||||
}
|
||||
NacosAggregateListener that = (NacosAggregateListener) o;
|
||||
return Objects.equals(notifyListener, that.notifyListener) && Objects.equals(serviceNames, that.serviceNames) && Objects.equals(serviceInstances, that.serviceInstances);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int hashCode() {
|
||||
return Objects.hash(notifyListener, serviceNames, serviceInstances);
|
||||
}
|
||||
}
|
||||
|
|
@ -51,6 +51,10 @@ public class NacosNamingServiceWrapper {
|
|||
namingService.subscribe(handleInnerSymbol(serviceName), group, eventListener);
|
||||
}
|
||||
|
||||
public void unsubscribe(String serviceName, String group, EventListener eventListener) throws NacosException {
|
||||
namingService.unsubscribe(handleInnerSymbol(serviceName), group, eventListener);
|
||||
}
|
||||
|
||||
public List<Instance> getAllInstances(String serviceName, String group) throws NacosException {
|
||||
return namingService.getAllInstances(handleInnerSymbol(serviceName), group);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -39,7 +39,6 @@ import com.alibaba.nacos.api.naming.listener.EventListener;
|
|||
import com.alibaba.nacos.api.naming.listener.NamingEvent;
|
||||
import com.alibaba.nacos.api.naming.pojo.Instance;
|
||||
import com.alibaba.nacos.api.naming.pojo.ListView;
|
||||
import com.alibaba.nacos.shaded.com.google.common.collect.Lists;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
|
|
@ -52,7 +51,6 @@ 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.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
|
@ -134,7 +132,9 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
*/
|
||||
private volatile ScheduledExecutorService scheduledExecutorService;
|
||||
|
||||
private final ConcurrentMap<URL, ConcurrentMap<NotifyListener, ConcurrentMap<String, EventListener>>> nacosListeners = new ConcurrentHashMap<>();
|
||||
private final Map<URL, Map<NotifyListener, NacosAggregateListener>> originToAggregateListener = new ConcurrentHashMap<>();
|
||||
|
||||
private final Map<URL, Map<NacosAggregateListener, Map<String, EventListener>>> nacosListeners = new ConcurrentHashMap<>();
|
||||
|
||||
public NacosRegistry(URL url, NacosNamingServiceWrapper namingService) {
|
||||
super(url);
|
||||
|
|
@ -203,7 +203,10 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
|
||||
@Override
|
||||
public void doSubscribe(final URL url, final NotifyListener listener) {
|
||||
Set<String> serviceNames = getServiceNames(url, listener);
|
||||
NacosAggregateListener nacosAggregateListener = new NacosAggregateListener(listener);
|
||||
originToAggregateListener.computeIfAbsent(url, k -> new ConcurrentHashMap<>()).put(listener, nacosAggregateListener);
|
||||
|
||||
Set<String> serviceNames = getServiceNames(url, nacosAggregateListener);
|
||||
|
||||
//Set corresponding serviceNames for easy search later
|
||||
if (isServiceNamesWithCompatibleMode(url)) {
|
||||
|
|
@ -211,14 +214,12 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
NacosInstanceManageUtil.setCorrespondingServiceNames(serviceName, serviceNames);
|
||||
}
|
||||
}
|
||||
|
||||
doSubscribe(url, listener, serviceNames);
|
||||
doSubscribe(url, nacosAggregateListener, serviceNames);
|
||||
}
|
||||
|
||||
private void doSubscribe(final URL url, final NotifyListener listener, final Set<String> serviceNames) {
|
||||
private void doSubscribe(final URL url, final NacosAggregateListener listener, final Set<String> serviceNames) {
|
||||
try {
|
||||
if (isServiceNamesWithCompatibleMode(url)) {
|
||||
List<Instance> allCorrespondingInstanceList = Lists.newArrayList();
|
||||
|
||||
/**
|
||||
* Get all instances with serviceNames to avoid instance overwrite and but with empty instance mentioned
|
||||
|
|
@ -233,9 +234,8 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
List<Instance> instances = namingService.getAllInstances(serviceName,
|
||||
getUrl().getGroup(Constants.DEFAULT_GROUP));
|
||||
NacosInstanceManageUtil.initOrRefreshServiceInstanceList(serviceName, instances);
|
||||
allCorrespondingInstanceList.addAll(instances);
|
||||
notifySubscriber(url, serviceName, listener, instances);
|
||||
}
|
||||
notifySubscriber(url, listener, allCorrespondingInstanceList);
|
||||
for (String serviceName : serviceNames) {
|
||||
subscribeEventListener(serviceName, url, listener);
|
||||
}
|
||||
|
|
@ -251,7 +251,7 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
}
|
||||
URL subscriberURL = url.setPath(serviceInterface).addParameters(INTERFACE_KEY, serviceInterface,
|
||||
CHECK_KEY, String.valueOf(false));
|
||||
notifySubscriber(subscriberURL, listener, instances);
|
||||
notifySubscriber(subscriberURL, serviceName, listener, instances);
|
||||
subscribeEventListener(serviceName, subscriberURL, listener);
|
||||
}
|
||||
}
|
||||
|
|
@ -275,6 +275,26 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
public void doUnsubscribe(URL url, NotifyListener listener) {
|
||||
if (isAdminProtocol(url)) {
|
||||
shutdownServiceNamesLookup();
|
||||
} else {
|
||||
Map<NotifyListener, NacosAggregateListener> listenerMap = originToAggregateListener.get(url);
|
||||
NacosAggregateListener nacosAggregateListener = listenerMap.remove(listener);
|
||||
if (nacosAggregateListener != null) {
|
||||
Set<String> serviceNames = getServiceNames(url, nacosAggregateListener);
|
||||
try {
|
||||
doUnsubscribe(url, nacosAggregateListener, serviceNames);
|
||||
} catch (NacosException e) {
|
||||
logger.error("Failed to unsubscribe " + url + " to nacos " + getUrl() + ", cause: " + e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
if (listenerMap.isEmpty()) {
|
||||
originToAggregateListener.remove(url);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void doUnsubscribe(final URL url, final NacosAggregateListener nacosAggregateListener, final Set<String> serviceNames) throws NacosException {
|
||||
for (String serviceName : serviceNames) {
|
||||
unsubscribeEventListener(serviceName, url, nacosAggregateListener);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -291,7 +311,7 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
* @param listener {@link NotifyListener}
|
||||
* @return non-null
|
||||
*/
|
||||
private Set<String> getServiceNames(URL url, NotifyListener listener) {
|
||||
private Set<String> getServiceNames(URL url, NacosAggregateListener listener) {
|
||||
if (isAdminProtocol(url)) {
|
||||
scheduleServiceNamesLookup(url, listener);
|
||||
return getServiceNamesForOps(url);
|
||||
|
|
@ -377,7 +397,7 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
return ADMIN_PROTOCOL.equals(url.getProtocol());
|
||||
}
|
||||
|
||||
private void scheduleServiceNamesLookup(final URL url, final NotifyListener listener) {
|
||||
private void scheduleServiceNamesLookup(final URL url, final NacosAggregateListener listener) {
|
||||
if (scheduledExecutorService == null) {
|
||||
scheduledExecutorService = Executors.newSingleThreadScheduledExecutor();
|
||||
scheduledExecutorService.scheduleAtFixedRate(() -> {
|
||||
|
|
@ -529,12 +549,12 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
return urls;
|
||||
}
|
||||
|
||||
private void subscribeEventListener(String serviceName, final URL url, final NotifyListener listener)
|
||||
private void subscribeEventListener(String serviceName, final URL url, final NacosAggregateListener listener)
|
||||
throws NacosException {
|
||||
ConcurrentMap<NotifyListener, ConcurrentMap<String, EventListener>> listeners = nacosListeners.computeIfAbsent(url,
|
||||
Map<NacosAggregateListener, Map<String, EventListener>> listeners = nacosListeners.computeIfAbsent(url,
|
||||
k -> new ConcurrentHashMap<>());
|
||||
|
||||
ConcurrentMap<String, EventListener> eventListeners = listeners.computeIfAbsent(listener,
|
||||
Map<String, EventListener> eventListeners = listeners.computeIfAbsent(listener,
|
||||
k -> new ConcurrentHashMap<>());
|
||||
|
||||
EventListener eventListener = eventListeners.computeIfAbsent(serviceName,
|
||||
|
|
@ -545,6 +565,30 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
eventListener);
|
||||
}
|
||||
|
||||
private void unsubscribeEventListener(String serviceName, final URL url, final NacosAggregateListener listener) throws NacosException {
|
||||
Map<NacosAggregateListener, Map<String, EventListener>> listenerToServiceEvent = nacosListeners.get(url);
|
||||
if (listenerToServiceEvent == null) {
|
||||
return;
|
||||
}
|
||||
Map<String, EventListener> serviceToEventMap = listenerToServiceEvent.get(listener);
|
||||
if (serviceToEventMap == null) {
|
||||
return;
|
||||
}
|
||||
EventListener eventListener = serviceToEventMap.remove(serviceName);
|
||||
if (eventListener == null) {
|
||||
return;
|
||||
}
|
||||
namingService.unsubscribe(serviceName,
|
||||
getUrl().getParameter(GROUP_KEY, Constants.DEFAULT_GROUP),
|
||||
eventListener);
|
||||
if (serviceToEventMap.isEmpty()) {
|
||||
listenerToServiceEvent.remove(listener);
|
||||
}
|
||||
if (listenerToServiceEvent.isEmpty()) {
|
||||
nacosListeners.remove(url);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Notify the Enabled {@link Instance instances} to subscriber.
|
||||
*
|
||||
|
|
@ -552,14 +596,14 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
* @param listener {@link NotifyListener}
|
||||
* @param instances all {@link Instance instances}
|
||||
*/
|
||||
private void notifySubscriber(URL url, NotifyListener listener, Collection<Instance> instances) {
|
||||
private void notifySubscriber(URL url, String serviceName, NacosAggregateListener listener, Collection<Instance> instances) {
|
||||
List<Instance> enabledInstances = new LinkedList<>(instances);
|
||||
if (enabledInstances.size() > 0) {
|
||||
// Instances
|
||||
filterEnabledInstances(enabledInstances);
|
||||
}
|
||||
List<URL> urls = toUrlWithEmpty(url, enabledInstances);
|
||||
NacosRegistry.this.notify(url, listener, urls);
|
||||
List<URL> aggregatedUrls = toUrlWithEmpty(url, listener.saveAndAggregateAllInstances(serviceName, enabledInstances));
|
||||
NacosRegistry.this.notify(url, listener.getNotifyListener(), aggregatedUrls);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -641,9 +685,9 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
|
||||
private final URL consumerUrl;
|
||||
|
||||
private final NotifyListener listener;
|
||||
private final NacosAggregateListener listener;
|
||||
|
||||
public RegistryChildListenerImpl(String serviceName, URL consumerUrl, NotifyListener listener) {
|
||||
public RegistryChildListenerImpl(String serviceName, URL consumerUrl, NacosAggregateListener listener) {
|
||||
this.serviceName = serviceName;
|
||||
this.consumerUrl = consumerUrl;
|
||||
this.listener = listener;
|
||||
|
|
@ -659,7 +703,7 @@ public class NacosRegistry extends FailbackRegistry {
|
|||
NacosInstanceManageUtil.initOrRefreshServiceInstanceList(serviceName, instances);
|
||||
instances = NacosInstanceManageUtil.getAllCorrespondingServiceInstanceList(serviceName);
|
||||
}
|
||||
NacosRegistry.this.notifySubscriber(consumerUrl, listener, instances);
|
||||
NacosRegistry.this.notifySubscriber(consumerUrl, serviceName, listener, instances);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue