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 <hyunwoojung@kakao.com>
This commit is contained in:
Hyunwoo Jung
2026-10-06 16:02:00 +02:00
committed by GitHub
parent 3709672494
commit 440139fec1
2 changed files with 105 additions and 8 deletions
@@ -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))
@@ -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<String> 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<String> 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<String> 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<String> 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 {