From 0dacabcd1dd6fa1c316b1488907bd7c1d3feb105 Mon Sep 17 00:00:00 2001 From: Ryan Baxter Date: Thu, 21 Jan 2021 11:35:05 -0500 Subject: [PATCH] Adding support for not ready addresses to informer based discovery client --- .../KubernetesInformerDiscoveryClient.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java index 4d7d77e1..920c872c 100644 --- a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java +++ b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java @@ -28,6 +28,7 @@ import io.kubernetes.client.extended.wait.Wait; import io.kubernetes.client.informer.SharedInformer; import io.kubernetes.client.informer.SharedInformerFactory; import io.kubernetes.client.informer.cache.Lister; +import io.kubernetes.client.openapi.models.V1EndpointAddress; import io.kubernetes.client.openapi.models.V1EndpointPort; import io.kubernetes.client.openapi.models.V1Endpoints; import io.kubernetes.client.openapi.models.V1Service; @@ -40,6 +41,7 @@ import org.springframework.cloud.client.discovery.DiscoveryClient; import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties; import org.springframework.cloud.kubernetes.commons.discovery.KubernetesServiceInstance; import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; public class KubernetesInformerDiscoveryClient implements DiscoveryClient, InitializingBean { @@ -130,7 +132,16 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi : subset.getPorts().stream() .filter(p -> this.properties.getPrimaryPortName().equalsIgnoreCase(p.getName())).findFirst() .orElseThrow(IllegalStateException::new); - return subset.getAddresses().stream() + List addresses = subset.getAddresses(); + if (this.properties.isIncludeNotReadyAddresses() + && !CollectionUtils.isEmpty(subset.getNotReadyAddresses())) { + if (addresses == null) { + addresses = new ArrayList<>(); + } + addresses.addAll(subset.getNotReadyAddresses()); + } + + return addresses.stream() .map(addr -> new KubernetesServiceInstance( addr.getTargetRef() != null ? addr.getTargetRef().getUid() : "", serviceId, addr.getIp(), port.getPort(), metadata, false));