fix migration: mark invokers as changed when refresh invoker

This commit is contained in:
ken.lj 2020-09-25 15:51:59 +08:00
parent c8e043a85b
commit e98a893387
3 changed files with 39 additions and 17 deletions

View File

@ -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<ServiceInstance> 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<ServiceInstance> 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));
}
/**

View File

@ -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<T> implements MigrationClusterInvoker<T> {
}
}
protected synchronized void discardServiceDiscoveryInvokerAddress(ClusterInvoker<?> serviceDiscoveryInvoker) {
protected synchronized void discardServiceDiscoveryInvokerAddress(ClusterInvoker<T> 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<Invoker<T>> 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<T> implements MigrationClusterInvoker<T> {
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<T> implements MigrationClusterInvoker<T> {
}
invoker = registryProtocol.getInvoker(cluster, registry, type, url);
} else {
((DynamicDirectory)invoker.getDirectory()).markInvokersChanged();
}
}
@ -331,7 +338,8 @@ public class MigrationInvoker<T> implements MigrationClusterInvoker<T> {
}
if (invoker != null) {
if (logger.isDebugEnabled()) {
logger.debug("Discarding interface addresses, total address size " + invoker.getDirectory().getAllInvokers().size());
List<Invoker<T>> invokers = invoker.getDirectory().getAllInvokers();
logger.debug("Discarding interface addresses, total address size " + (invokers == null ? 0 : invokers.size()));
}
//invoker.getDirectory().discordAddresses();
}

View File

@ -260,17 +260,27 @@ public abstract class DynamicDirectory<T> extends AbstractDirectory<T> 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();
}