From 440139fec1e61b275ea152da8293e7647ab9b917 Mon Sep 17 00:00:00 2001 From: Hyunwoo Jung Date: Tue, 6 Oct 2026 23:02:00 +0900 Subject: [PATCH] Fix Flux resubscription in reactive caching Fix an existing bug in processPutRequest() affecting `@Cacheable` and `@CachePut` since 6.1: resubscription after an error hangs while publish().refCount(2) waits for a second subscriber. Also fix an unreleased regression introduced by 26de340 (gh-37309) for `@CacheEvict`, which copied the same pattern into processCacheEvicts(). Create the shared Flux and cache subscriber for each subscription via Flux.defer(). Add regression tests combining reactive caching with `@Retryable`. See gh-37309 Closes gh-37403 Signed-off-by: Hyunwoo Jung --- .../cache/interceptor/CacheAspectSupport.java | 22 +++-- .../ReactiveRetryInterceptorTests.java | 91 +++++++++++++++++++ 2 files changed, 105 insertions(+), 8 deletions(-) diff --git a/spring-context/src/main/java/org/springframework/cache/interceptor/CacheAspectSupport.java b/spring-context/src/main/java/org/springframework/cache/interceptor/CacheAspectSupport.java index 37814657e3b..041956de272 100644 --- a/spring-context/src/main/java/org/springframework/cache/interceptor/CacheAspectSupport.java +++ b/spring-context/src/main/java/org/springframework/cache/interceptor/CacheAspectSupport.java @@ -1220,10 +1220,13 @@ public abstract class CacheAspectSupport extends AbstractCacheInvoker ReactiveAdapter adapter = (result != null ? this.registry.getAdapter(result.getClass()) : null); if (adapter != null) { if (adapter.isMultiValue()) { - Flux source = Flux.from(adapter.toPublisher(result)) - .publish().refCount(2); - source.subscribe(new CacheEvictListSubscriber(contexts)); - return adapter.fromPublisher(source); + Flux deferred = Flux.defer(() -> { + Flux source = Flux.from(adapter.toPublisher(result)) + .publish().refCount(2); + source.subscribe(new CacheEvictListSubscriber(contexts)); + return source; + }); + return adapter.fromPublisher(deferred); } else { return adapter.fromPublisher(Mono.from(adapter.toPublisher(result)) @@ -1287,10 +1290,13 @@ public abstract class CacheAspectSupport extends AbstractCacheInvoker ReactiveAdapter adapter = (result != null ? this.registry.getAdapter(result.getClass()) : null); if (adapter != null) { if (adapter.isMultiValue()) { - Flux source = Flux.from(adapter.toPublisher(result)) - .publish().refCount(2); - source.subscribe(new CachePutListSubscriber(request)); - return adapter.fromPublisher(source); + Flux deferred = Flux.defer(() -> { + Flux source = Flux.from(adapter.toPublisher(result)) + .publish().refCount(2); + source.subscribe(new CachePutListSubscriber(request)); + return source; + }); + return adapter.fromPublisher(deferred); } else { return adapter.fromPublisher(Mono.from(adapter.toPublisher(result)) diff --git a/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java b/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java index 50bf4428094..e4e7a8140d0 100644 --- a/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java +++ b/spring-context/src/test/java/org/springframework/resilience/ReactiveRetryInterceptorTests.java @@ -22,6 +22,7 @@ import java.nio.charset.MalformedInputException; import java.nio.file.AccessDeniedException; import java.nio.file.FileSystemException; import java.time.Duration; +import java.util.List; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; @@ -36,6 +37,13 @@ import org.springframework.aop.framework.AopProxyUtils; import org.springframework.aop.framework.ProxyFactory; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.cache.Cache; +import org.springframework.cache.CacheManager; +import org.springframework.cache.annotation.CacheEvict; +import org.springframework.cache.annotation.Cacheable; +import org.springframework.cache.annotation.EnableCaching; +import org.springframework.cache.concurrent.ConcurrentMapCacheManager; +import org.springframework.cache.interceptor.SimpleKey; import org.springframework.context.annotation.AnnotationConfigApplicationContext; import org.springframework.context.support.GenericApplicationContext; import org.springframework.resilience.annotation.EnableResilientMethods; @@ -189,6 +197,49 @@ class ReactiveRetryInterceptorTests { assertThat(target.counter).hasValue(2); } + @Test // gh-37403 + void withCacheableAnnotation() { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext(); + ctx.registerBeanDefinition("bean", new RootBeanDefinition(CacheableAnnotatedBean.class)); + ctx.registerBeanDefinition("cacheManager", new RootBeanDefinition(ConcurrentMapCacheManager.class)); + ctx.registerBeanDefinition("config", new RootBeanDefinition(EnablingConfigWithCaching.class)); + ctx.refresh(); + CacheableAnnotatedBean proxy = ctx.getBean(CacheableAnnotatedBean.class); + CacheableAnnotatedBean target = (CacheableAnnotatedBean) AopProxyUtils.getSingletonTarget(proxy); + + // Simulates a failure on the 1st subscription to the target Flux. + List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5)); + assertThat(result).containsExactly("a", "b"); + // Subscribed twice: the 2nd attempt re-subscribed to the Flux returned from + // the cache interceptor and reached the target Flux again instead of hanging. + assertThat(target.counter).hasValue(2); + + Cache cache = ctx.getBean(CacheManager.class).getCache("tests"); + assertThat(cache.get(SimpleKey.EMPTY).get()).isEqualTo(List.of("a", "b")); + } + + @Test // gh-37403 + void withCacheEvictAnnotation() { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext(); + ctx.registerBeanDefinition("bean", new RootBeanDefinition(CacheEvictAnnotatedBean.class)); + ctx.registerBeanDefinition("cacheManager", new RootBeanDefinition(ConcurrentMapCacheManager.class)); + ctx.registerBeanDefinition("config", new RootBeanDefinition(EnablingConfigWithCaching.class)); + ctx.refresh(); + CacheEvictAnnotatedBean proxy = ctx.getBean(CacheEvictAnnotatedBean.class); + CacheEvictAnnotatedBean target = (CacheEvictAnnotatedBean) AopProxyUtils.getSingletonTarget(proxy); + + Cache cache = ctx.getBean(CacheManager.class).getCache("tests"); + cache.put(SimpleKey.EMPTY, List.of("a", "b")); + + // Simulates a failure on the 1st subscription to the target Flux. + List result = proxy.retryOperation().collectList().block(Duration.ofSeconds(5)); + assertThat(result).containsExactly("a", "b"); + // Subscribed twice: the 2nd attempt re-subscribed to the Flux returned from + // the cache interceptor and reached the target Flux again instead of hanging. + assertThat(target.counter).hasValue(2); + assertThat(cache.get(SimpleKey.EMPTY)).isNull(); + } + @Test void withMethodRetryEventListener() throws Exception { AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext(); @@ -636,11 +687,51 @@ class ReactiveRetryInterceptorTests { } + static class CacheableAnnotatedBean { + + AtomicInteger counter = new AtomicInteger(); + + @Cacheable("tests") + @Retryable(maxRetries = 2, delay = 10) + public Flux retryOperation() { + return Flux.defer(() -> { + if (counter.incrementAndGet() == 1) { + return Flux.error(new IOException(counter.toString())); + } + return Flux.just("a", "b"); + }); + } + } + + + static class CacheEvictAnnotatedBean { + + AtomicInteger counter = new AtomicInteger(); + + @CacheEvict("tests") + @Retryable(maxRetries = 2, delay = 10) + public Flux retryOperation() { + return Flux.defer(() -> { + if (counter.incrementAndGet() == 1) { + return Flux.error(new IOException(counter.toString())); + } + return Flux.just("a", "b"); + }); + } + } + + @EnableResilientMethods static class EnablingConfig { } + @EnableCaching + @EnableResilientMethods + static class EnablingConfigWithCaching { + } + + // Bean classes for boundary testing static class MinimalRetryBean {