diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/route/CachingRouteLocator.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/route/CachingRouteLocator.java index d7e1a5a8..b010f7aa 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/route/CachingRouteLocator.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/route/CachingRouteLocator.java @@ -26,6 +26,7 @@ import org.apache.commons.logging.LogFactory; import reactor.cache.CacheFlux; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.core.publisher.Signal; import org.springframework.cloud.gateway.event.RefreshRoutesEvent; import org.springframework.cloud.gateway.event.RefreshRoutesResultEvent; @@ -90,20 +91,16 @@ public class CachingRouteLocator scopedRoutes.subscribe(scopedRoutesList -> { Flux.concat(Flux.fromIterable(scopedRoutesList), getNonScopedRoutes(event)) .sort(AnnotationAwareOrderComparator.INSTANCE).materialize().collect(Collectors.toList()) - .subscribe(signals -> { - cache.put(CACHE_KEY, signals); - applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this)); - }, this::handleRefreshError); + .subscribe(this::publishRefreshEvent, this::handleRefreshError); }, this::handleRefreshError); } else { final Mono> allRoutes = fetch().collect(Collectors.toList()); - allRoutes.subscribe(list -> Flux.fromIterable(list).materialize().collect(Collectors.toList()) - .subscribe(signals -> { - cache.put(CACHE_KEY, signals); - applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this)); - }, this::handleRefreshError), this::handleRefreshError); + allRoutes.subscribe( + list -> Flux.fromIterable(list).materialize().collect(Collectors.toList()) + .subscribe(this::publishRefreshEvent, this::handleRefreshError), + this::handleRefreshError); } } catch (Throwable e) { @@ -111,6 +108,11 @@ public class CachingRouteLocator } } + private void publishRefreshEvent(List> signals) { + cache.put(CACHE_KEY, signals); + applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this)); + } + private Flux getNonScopedRoutes(RefreshRoutesEvent scopedEvent) { return this.getRoutes() .filter(route -> !RouteLocator.matchMetadata(route.getMetadata(), scopedEvent.getMetadata()));