From cc52e6cefc9cf3632cacc90da4941b6bab88bacd Mon Sep 17 00:00:00 2001 From: "ken.lj" Date: Mon, 9 Nov 2020 22:22:45 +0800 Subject: [PATCH] fix instance listener metadata update --- .../ServiceInstancesChangedListener.java | 131 +++++++++--------- 1 file changed, 69 insertions(+), 62 deletions(-) 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 e00e98a013..ba05262d37 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 @@ -101,24 +101,12 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener * @param event {@link ServiceInstancesChangedEvent} */ public synchronized void onEvent(ServiceInstancesChangedEvent event) { - String appName = event.getServiceName(); - List appInstances = event.getServiceInstances(); - if (event instanceof RetryServiceInstancesChangedEvent) { - RetryServiceInstancesChangedEvent retryEvent = (RetryServiceInstancesChangedEvent) event; - logger.warn("Received address refresh retry event, " + retryEvent.getFailureRecordTime()); - if (retryEvent.getFailureRecordTime() < lastRefreshTime) { - logger.warn("Ignore retry event, event time: " + retryEvent.getFailureRecordTime() + ", last refresh time: " + lastRefreshTime); - return; - } - logger.warn("Retrying address notification..."); - } else { - logger.info("Received instance notification, serviceName: " + appName + ", instances: " + appInstances.size()); - allInstances.put(appName, appInstances); - lastRefreshTime = System.currentTimeMillis(); + if (this.isRetryAndExpired(event)) { + return; } if (logger.isDebugEnabled()) { - logger.debug(appInstances.toString()); + logger.debug(event.getServiceInstances().toString()); } Map> revisionToInstances = new HashMap<>(); @@ -151,15 +139,11 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener scheduler.submit(new AddressRefreshRetryTask(retryPermission)); logger.warn("Address refresh try task submitted."); } - logger.warn("Address refresh failed because of Metadata Server failure, wait for retry or new refresh event."); + logger.warn("Address refresh failed because of Metadata Server failure, wait for retry or new address refresh event."); this.revisionToMetadata = newRevisionToMetadata; return; } - if (revisionToMetadata.size() != 0) { - logger.info("Revisions removed: " + revisionToMetadata.keySet()); - revisionToMetadata.clear(); - } this.revisionToMetadata = newRevisionToMetadata; localServiceToRevisions.forEach((serviceKey, revisions) -> { @@ -182,6 +166,67 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener this.notifyAddressChanged(); } + public void addListener(String serviceKey, NotifyListener listener) { + this.listeners.put(serviceKey, listener); + } + + public void removeListener(String serviceKey) { + listeners.remove(serviceKey); + if (listeners.isEmpty()) { + serviceDiscovery.removeServiceInstancesChangedListener(this); + } + } + + public List getUrls(String serviceKey) { + return toUrlsWithEmpty(serviceUrls.get(serviceKey)); + } + + /** + * Get the correlative service name + * + * @return the correlative service name + */ + public final Set getServiceNames() { + return serviceNames; + } + + public void setUrl(URL url) { + this.url = url; + } + + public URL getUrl() { + return url; + } + + /** + * @param event {@link ServiceInstancesChangedEvent event} + * @return If service name matches, return true, or false + */ + public final boolean accept(ServiceInstancesChangedEvent event) { + return serviceNames.contains(event.getServiceName()); + } + + + private boolean isRetryAndExpired(ServiceInstancesChangedEvent event) { + String appName = event.getServiceName(); + List appInstances = event.getServiceInstances(); + + if (event instanceof RetryServiceInstancesChangedEvent) { + RetryServiceInstancesChangedEvent retryEvent = (RetryServiceInstancesChangedEvent) event; + logger.warn("Received address refresh retry event, " + retryEvent.getFailureRecordTime()); + if (retryEvent.getFailureRecordTime() < lastRefreshTime) { + logger.warn("Ignore retry event, event time: " + retryEvent.getFailureRecordTime() + ", last refresh time: " + lastRefreshTime); + return true; + } + logger.warn("Retrying address notification..."); + } else { + logger.info("Received instance notification, serviceName: " + appName + ", instances: " + appInstances.size()); + allInstances.put(appName, appInstances); + lastRefreshTime = System.currentTimeMillis(); + } + return false; + } + private boolean hasEmptyMetadata(Map revisionToMetadata) { if (revisionToMetadata == null) { return false; @@ -197,13 +242,14 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener } private MetadataInfo getRemoteMetadata(ServiceInstance instance, String revision, Map> localServiceToRevisions, List subInstances) { - MetadataInfo metadata = revisionToMetadata.remove(revision); + MetadataInfo metadata = revisionToMetadata.get(revision); if (metadata == null) { if (failureCounter.get() < 3 || (System.currentTimeMillis() - lastFailureTime > 5000)) { metadata = getMetadataInfo(instance); if (metadata != null) { logger.info("MetadataInfo for instance " + instance.getAddress() + "?revision=" + revision + " is " + metadata); failureCounter.set(0); + revisionToMetadata.put(revision, metadata); parseMetadata(revision, metadata, localServiceToRevisions); } else { logger.error("Failed to get MetadataInfo for instance " + instance.getAddress() + "?revision=" + revision @@ -212,7 +258,8 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener failureCounter.incrementAndGet(); } } - } else if (subInstances.size() > 1) {// check if metadata info parsed + } else if (subInstances.size() < 1) { + // "subInstances.size() >= 1" means metadata of this revision has been parsed, ignore parseMetadata(revision, metadata, localServiceToRevisions); } return metadata; @@ -268,46 +315,6 @@ public class ServiceInstancesChangedListener implements ConditionalEventListener return urls; } - public void addListener(String serviceKey, NotifyListener listener) { - this.listeners.put(serviceKey, listener); - } - - public void removeListener(String serviceKey) { - listeners.remove(serviceKey); - if (listeners.isEmpty()) { - serviceDiscovery.removeServiceInstancesChangedListener(this); - } - } - - public List getUrls(String serviceKey) { - return toUrlsWithEmpty(serviceUrls.get(serviceKey)); - } - - /** - * Get the correlative service name - * - * @return the correlative service name - */ - public final Set getServiceNames() { - return serviceNames; - } - - public void setUrl(URL url) { - this.url = url; - } - - public URL getUrl() { - return url; - } - - /** - * @param event {@link ServiceInstancesChangedEvent event} - * @return If service name matches, return true, or false - */ - public final boolean accept(ServiceInstancesChangedEvent event) { - return serviceNames.contains(event.getServiceName()); - } - @Override public boolean equals(Object o) { if (this == o) return true;