You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 18:26:59