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 b010f7aa..52b23cd3 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 @@ -89,18 +89,13 @@ public class CachingRouteLocator .onErrorResume(s -> Mono.just(List.of())); scopedRoutes.subscribe(scopedRoutesList -> { - Flux.concat(Flux.fromIterable(scopedRoutesList), getNonScopedRoutes(event)) - .sort(AnnotationAwareOrderComparator.INSTANCE).materialize().collect(Collectors.toList()) - .subscribe(this::publishRefreshEvent, this::handleRefreshError); + updateCache(Flux.concat(Flux.fromIterable(scopedRoutesList), getNonScopedRoutes(event)) + .sort(AnnotationAwareOrderComparator.INSTANCE)); }, this::handleRefreshError); } else { final Mono> allRoutes = fetch().collect(Collectors.toList()); - - allRoutes.subscribe( - list -> Flux.fromIterable(list).materialize().collect(Collectors.toList()) - .subscribe(this::publishRefreshEvent, this::handleRefreshError), - this::handleRefreshError); + allRoutes.subscribe(list -> updateCache(Flux.fromIterable(list)), this::handleRefreshError); } } catch (Throwable e) { @@ -108,6 +103,11 @@ public class CachingRouteLocator } } + private synchronized void updateCache(Flux routes) { + routes.materialize().collect(Collectors.toList()).subscribe(this::publishRefreshEvent, + this::handleRefreshError); + } + private void publishRefreshEvent(List> signals) { cache.put(CACHE_KEY, signals); applicationEventPublisher.publishEvent(new RefreshRoutesResultEvent(this));