diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocator.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocator.java index 80f8f44b..8109fc18 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocator.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocator.java @@ -17,6 +17,7 @@ package org.springframework.cloud.gateway.discovery; import java.net.URI; +import java.util.List; import java.util.Map; import java.util.function.Predicate; @@ -27,6 +28,7 @@ import reactor.core.scheduler.Schedulers; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.gateway.filter.FilterDefinition; import org.springframework.cloud.gateway.handler.predicate.PredicateDefinition; import org.springframework.cloud.gateway.route.RouteDefinition; @@ -49,23 +51,45 @@ public class DiscoveryClientRouteDefinitionLocator implements RouteDefinitionLoc private static final Log log = LogFactory .getLog(DiscoveryClientRouteDefinitionLocator.class); - private final DiscoveryClient discoveryClient; - private final DiscoveryLocatorProperties properties; private final String routeIdPrefix; private final SimpleEvaluationContext evalCtxt; + private Flux> serviceInstances; + + /** + * Kept for backwards compatibility. You should use the reactive discovery client. + * @param discoveryClient the blocking discovery client + * @param properties the configuration properties + * @deprecated kept for backwards compatibility + */ + @Deprecated public DiscoveryClientRouteDefinitionLocator(DiscoveryClient discoveryClient, DiscoveryLocatorProperties properties) { - this.discoveryClient = discoveryClient; + this(discoveryClient.getClass().getSimpleName(), properties); + serviceInstances = Flux + .defer(() -> Flux.fromIterable(discoveryClient.getServices())) + .map(discoveryClient::getInstances) + .subscribeOn(Schedulers.boundedElastic()); + } + + public DiscoveryClientRouteDefinitionLocator(ReactiveDiscoveryClient discoveryClient, + DiscoveryLocatorProperties properties) { + this(discoveryClient.getClass().getSimpleName(), properties); + serviceInstances = discoveryClient.getServices() + .flatMap(service -> discoveryClient.getInstances(service).collectList()); + } + + private DiscoveryClientRouteDefinitionLocator(String discoveryClientName, + DiscoveryLocatorProperties properties) { this.properties = properties; if (StringUtils.hasText(properties.getRouteIdPrefix())) { - this.routeIdPrefix = properties.getRouteIdPrefix(); + routeIdPrefix = properties.getRouteIdPrefix(); } else { - this.routeIdPrefix = this.discoveryClient.getClass().getSimpleName() + "_"; + routeIdPrefix = discoveryClientName + "_"; } evalCtxt = SimpleEvaluationContext.forReadOnlyDataBinding().withInstanceMethods() .build(); @@ -94,9 +118,7 @@ public class DiscoveryClientRouteDefinitionLocator implements RouteDefinitionLoc }; } - return Flux.defer(() -> Flux.fromIterable(discoveryClient.getServices())) - .subscribeOn(Schedulers.elastic()).map(discoveryClient::getInstances) - .filter(instances -> !instances.isEmpty()) + return serviceInstances.filter(instances -> !instances.isEmpty()) .map(instances -> instances.get(0)).filter(includePredicate) .map(instance -> { String serviceId = instance.getServiceId(); diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfiguration.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfiguration.java index ee9dacd0..60e92b5d 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfiguration.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfiguration.java @@ -25,7 +25,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.client.discovery.composite.CompositeDiscoveryClientAutoConfiguration; import org.springframework.cloud.gateway.config.GatewayAutoConfiguration; import org.springframework.cloud.gateway.filter.FilterDefinition; @@ -49,7 +49,7 @@ import static org.springframework.cloud.gateway.support.NameUtils.normalizeRoute @ConditionalOnProperty(name = "spring.cloud.gateway.enabled", matchIfMissing = true) @AutoConfigureBefore(GatewayAutoConfiguration.class) @AutoConfigureAfter(CompositeDiscoveryClientAutoConfiguration.class) -@ConditionalOnClass({ DispatcherHandler.class, DiscoveryClient.class }) +@ConditionalOnClass({ DispatcherHandler.class, ReactiveDiscoveryClient.class }) @EnableConfigurationProperties public class GatewayDiscoveryClientAutoConfiguration { @@ -81,10 +81,11 @@ public class GatewayDiscoveryClientAutoConfiguration { } @Bean - @ConditionalOnBean(DiscoveryClient.class) + @ConditionalOnBean(ReactiveDiscoveryClient.class) @ConditionalOnProperty(name = "spring.cloud.gateway.discovery.locator.enabled") public DiscoveryClientRouteDefinitionLocator discoveryClientRouteDefinitionLocator( - DiscoveryClient discoveryClient, DiscoveryLocatorProperties properties) { + ReactiveDiscoveryClient discoveryClient, + DiscoveryLocatorProperties properties) { return new DiscoveryClientRouteDefinitionLocator(discoveryClient, properties); } diff --git a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorIntegrationTests.java b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorIntegrationTests.java index 87ae2182..1d1aaf00 100644 --- a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorIntegrationTests.java +++ b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorIntegrationTests.java @@ -16,13 +16,12 @@ package org.springframework.cloud.gateway.discovery; -import java.util.Arrays; -import java.util.Collections; import java.util.List; import java.util.concurrent.atomic.AtomicBoolean; import org.junit.Test; import org.junit.runner.RunWith; +import reactor.core.publisher.Flux; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringBootConfiguration; @@ -30,7 +29,7 @@ import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.client.DefaultServiceInstance; import org.springframework.cloud.client.ServiceInstance; -import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.client.discovery.event.HeartbeatEvent; import org.springframework.cloud.gateway.route.Route; import org.springframework.cloud.gateway.route.RouteLocator; @@ -84,7 +83,7 @@ public class DiscoveryClientRouteDefinitionLocatorIntegrationTests { } - private static class TestDiscoveryClient implements DiscoveryClient { + private static class TestDiscoveryClient implements ReactiveDiscoveryClient { AtomicBoolean single = new AtomicBoolean(true); @@ -104,22 +103,22 @@ public class DiscoveryClientRouteDefinitionLocatorIntegrationTests { } @Override - public List getInstances(String serviceId) { + public Flux getInstances(String serviceId) { if (serviceId.equals("service1")) { - return Collections.singletonList(instance1); + return Flux.just(instance1); } if (serviceId.equals("service2")) { - return Collections.singletonList(instance2); + return Flux.just(instance2); } - return Collections.emptyList(); + return Flux.empty(); } @Override - public List getServices() { + public Flux getServices() { if (single.get()) { - return Collections.singletonList("service1"); + return Flux.just("service1"); } - return Arrays.asList("service1", "service2"); + return Flux.just("service1", "service2"); } } diff --git a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorTests.java b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorTests.java index 52ba7409..d44ac252 100644 --- a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorTests.java +++ b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/DiscoveryClientRouteDefinitionLocatorTests.java @@ -16,20 +16,20 @@ package org.springframework.cloud.gateway.discovery; -import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Map; import org.junit.Test; import org.junit.runner.RunWith; +import reactor.core.publisher.Flux; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringBootConfiguration; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.client.DefaultServiceInstance; -import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.gateway.filter.FilterDefinition; import org.springframework.cloud.gateway.handler.predicate.PredicateDefinition; import org.springframework.cloud.gateway.route.RouteDefinition; @@ -95,22 +95,22 @@ public class DiscoveryClientRouteDefinitionLocatorTests { protected static class Config { @Bean - DiscoveryClient discoveryClient() { - DiscoveryClient discoveryClient = mock(DiscoveryClient.class); + ReactiveDiscoveryClient discoveryClient() { + ReactiveDiscoveryClient discoveryClient = mock(ReactiveDiscoveryClient.class); when(discoveryClient.getServices()) - .thenReturn(Arrays.asList("SERVICE1", "Service2")); + .thenReturn(Flux.just("SERVICE1", "Service2")); whenInstance(discoveryClient, "SERVICE1", Collections.singletonMap("edge", "true")); whenInstance(discoveryClient, "Service2", Collections.emptyMap()); return discoveryClient; } - private void whenInstance(DiscoveryClient discoveryClient, String serviceId, - Map metadata) { + private void whenInstance(ReactiveDiscoveryClient discoveryClient, + String serviceId, Map metadata) { DefaultServiceInstance instance1 = new DefaultServiceInstance(serviceId, "localhost", 8001, false, metadata); when(discoveryClient.getInstances(serviceId)) - .thenReturn(Collections.singletonList(instance1)); + .thenReturn(Flux.just(instance1)); } } diff --git a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfigurationTests.java b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfigurationTests.java index 33f45b95..c0151208 100644 --- a/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfigurationTests.java +++ b/spring-cloud-gateway-core/src/test/java/org/springframework/cloud/gateway/discovery/GatewayDiscoveryClientAutoConfigurationTests.java @@ -19,18 +19,20 @@ package org.springframework.cloud.gateway.discovery; import org.junit.Test; import org.junit.experimental.runners.Enclosed; import org.junit.runner.RunWith; +import reactor.core.publisher.Flux; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringBootConfiguration; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.gateway.config.LoadBalancerProperties; import org.springframework.context.annotation.Bean; import org.springframework.test.context.junit4.SpringRunner; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; @RunWith(Enclosed.class) public class GatewayDiscoveryClientAutoConfigurationTests { @@ -80,8 +82,10 @@ public class GatewayDiscoveryClientAutoConfigurationTests { protected static class Config { @Bean - DiscoveryClient discoveryClient() { - return mock(DiscoveryClient.class); + ReactiveDiscoveryClient discoveryClient() { + ReactiveDiscoveryClient discoveryClient = mock(ReactiveDiscoveryClient.class); + when(discoveryClient.getServices()).thenReturn(Flux.empty()); + return discoveryClient; } }