Webflux集成Reactive Redis缓存序列化异常求助
问题描述
已在CacheConfig中配置ReactiveRedisConnectionFactory并成功连接Redis,但使用Spring标准@Cacheable注解缓存Flux<Place>数据时失败,抛出如下异常:
Caused by: java.lang.IllegalArgumentException: DefaultSerializer requires a Serializable payload but received an object of type [reactor.core.publisher.MonoFlatMapMany] at org.springframework.core.serializer.DefaultSerializer.serialize(DefaultSerializer.java:43) ~[spring-core-6.0.7.jar:6.0.7] at org.springframework.core.serializer.Serializer.serializeToByteArray(Serializer.java:56) ~[spring-core-6.0.7.jar:6.0.7] at org.springframework.core.serializer.support.SerializingConverter.convert(SerializingConverter.java:60) ~[spring-core-6.0.7.jar:6.0.7]
问题源于无法直接序列化Flux/Mono对象,尝试自定义ReactiveRedisSerializer处理,但必须调用Flux的block()方法获取内容转为byte[],这不符合非阻塞设计要求,想了解是否有官方支持的开箱即用非阻塞解决方案。
原缓存配置代码:
@EnableCaching @Configuration @RequiredArgsConstructor public class CacheConfig implements CachingConfigurer { private final ObjectMapper objectMapper; @Bean public LettuceConnectionFactory redisConnectionFactory() { LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder() .commandTimeout(Duration.ofSeconds(2)) .shutdownTimeout(Duration.ZERO) .build(); RedisStandaloneConfiguration serverConfig = new RedisStandaloneConfiguration("localhost", 6379); return new LettuceConnectionFactory(serverConfig, clientConfig); } @Bean public CacheManager cacheManager() { RedisCacheConfiguration cacheConfig = RedisCacheConfiguration.defaultCacheConfig() .entryTtl(Duration.ofMinutes(5)) .serializeValuesWith(RedisSerializationContext.SerializationPair.fromSerializer(new ReactiveRedisSerializer<>(objectMapper))); return RedisCacheManager.builder(redisConnectionFactory()) .cacheDefaults(cacheConfig) .build(); } @Override public CacheResolver cacheResolver() { return new SimpleCacheResolver(cacheManager()); } }
自定义序列化器(存在阻塞问题):
@Slf4j @RequiredArgsConstructor public class ReactiveRedisSerializer<T> implements RedisSerializer<T> { private final ObjectMapper objectMapper; @Override public byte[] serialize(T t) throws SerializationException { if (t == null) { return null; } if (t instanceof Mono<?> mono) { try { return objectMapper.writeValueAsBytes(mono.toFuture().get()); } catch (JsonProcessingException | InterruptedException | ExecutionException e) { throw new SerializationException("Failed to serialize Mono", e); } } if (t instanceof Flux<?> flux) { try { List<?> list = flux.collectList().block(); // How to avoid doing this? return objectMapper.writeValueAsBytes(list); } catch (JsonProcessingException e) { throw new SerializationException("Failed to serialize Flux", e); } } try { return objectMapper.writeValueAsString(t).getBytes(StandardCharsets.UTF_8); } catch (JsonProcessingException e) { throw new SerializationException("Failed to serialize object", e); } } @Override public T deserialize(byte[] bytes) throws SerializationException { if (bytes == null) { return null; } try { return (T) objectMapper.readValue(bytes, Object.class); } catch (Exception e) { throw new SerializationException("Failed to deserialize object", e); } } }
业务缓存方法:
@Cacheable(cacheNames = "placesCache") public Flux<Place> getPlaces(final PlaceRequest placeRequest) { final String query = "%s in %s".formatted(placeRequest.vendorType(), placeRequest.city()); return googleMapsService .textSearch(query) .flatMapMany(response -> placeMapper.toPlaces(response) .map(place -> place.withQuery(query).withVendorType(placeRequest.vendorType()))); }
解决方案
1. 替换为响应式缓存管理器
普通的RedisCacheManager是为同步场景设计的,无法适配Flux/Mono这类响应式类型的非阻塞缓存需求。必须改用Spring Data Redis提供的ReactiveRedisCacheManager,它完全支持响应式流的非阻塞缓存操作。
2. 使用官方序列化器替代自定义实现
无需手动编写带block()的序列化器,Spring Data Redis提供了Jackson2JsonRedisSerializer,可以直接集成到响应式缓存配置中,实现非阻塞的序列化/反序列化。
修改后的缓存配置
@EnableCaching @Configuration @RequiredArgsConstructor public class CacheConfig implements CachingConfigurer { private final ObjectMapper objectMapper; @Bean public LettuceConnectionFactory redisConnectionFactory() { LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder() .commandTimeout(Duration.ofSeconds(2)) .shutdownTimeout(Duration.ZERO) .build(); RedisStandaloneConfiguration serverConfig = new RedisStandaloneConfiguration("localhost", 6379); return new LettuceConnectionFactory(serverConfig, clientConfig); } // 配置响应式缓存管理器,适配Flux/Mono @Bean public ReactiveCacheManager reactiveCacheManager() { RedisCacheConfiguration cacheConfig = RedisCacheConfiguration.defaultCacheConfig() .entryTtl(Duration.ofMinutes(5)) .serializeValuesWith(RedisSerializationContext.SerializationPair.fromSerializer( new Jackson2JsonRedisSerializer<>(objectMapper, Object.class) )); return RedisCacheManager.builder(redisConnectionFactory()) .cacheDefaults(cacheConfig) .buildReactiveCacheManager(); } // 保留同步缓存管理器(如果有同步场景需要) @Override public CacheManager cacheManager() { RedisCacheConfiguration cacheConfig = RedisCacheConfiguration.defaultCacheConfig() .entryTtl(Duration.ofMinutes(5)) .serializeValuesWith(RedisSerializationContext.SerializationPair.fromSerializer( new Jackson2JsonRedisSerializer<>(objectMapper, Object.class) )); return RedisCacheManager.builder(redisConnectionFactory()) .cacheDefaults(cacheConfig) .build(); } @Override public CacheResolver cacheResolver() { return new SimpleCacheResolver(cacheManager()); } }
3. 无需修改业务方法的缓存注解
Spring会自动识别返回值为Flux/Mono的@Cacheable方法,并使用ReactiveCacheManager进行非阻塞缓存处理,不需要手动订阅或阻塞流。
关键原理说明
ReactiveRedisCacheManager会订阅响应式流(Flux/Mono),获取实际数据后再进行序列化缓存,全程非阻塞。- 官方序列化器
Jackson2JsonRedisSerializer直接处理实际业务对象(比如List<Place>),而非Flux/Mono容器,避免了序列化容器对象的问题。
额外注意事项
- 确保Spring Data Redis版本与Spring Framework版本兼容(Spring 6.x对应Spring Data Redis 3.x及以上)。
- 如果需要精确反序列化类型,可以将
Jackson2JsonRedisSerializer的泛型指定为List<Place>,或者通过objectMapper配置类型信息,防止反序列化时类型丢失。
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

