如何将Mono<Flux<T>>转换为Flux<T>(Spring Webflux+R2DBC场景)
解决Mono<Flux>转Flux的问题
你的问题出在collectList()之后调用map(brandRepo::saveAll):因为saveAll(List<Brand>)返回Flux<Brand>,所以Mono<List<Brand>>经过map转换后就变成了Mono<Flux<Brand>>。这里提供两种解决思路:
方案一:使用flatMapMany替换map
flatMapMany操作符可以将Mono内部包裹的Flux直接展开为顶层的Flux类型,适配你的需求:
public Flux<Brand> refreshBrands() { return Flux.fromIterable(commonService.parseProduct())// 获取产品列表 .flatMap(this::fetchBrandsWithWebClient)// 转换为品牌数据 .filter(Brand::isNotEmptyBrand)// 过滤无效品牌 .distinct()// 去重 .collectList()// 收集为List .flatMapMany(brandRepo::saveAll);// 保存并展开Flux }
方案二:去掉collectList,直接流式保存(推荐)
R2DBC的Reactive Repository提供了接受Publisher<T>(包括Flux)的saveAll重载方法,完全不需要先把数据收集到List中。这种方式符合响应式流式处理的设计,避免内存过载:
public Flux<Brand> refreshBrands() { return Flux.fromIterable(commonService.parseProduct()) .flatMap(this::fetchBrandsWithWebClient) .filter(Brand::isNotEmptyBrand) .distinct() .as(brandRepo::saveAll); // 直接将Flux传入saveAll } // 更直观的写法: public Flux<Brand> refreshBrands() { Flux<Brand> brandStream = Flux.fromIterable(commonService.parseProduct()) .flatMap(this::fetchBrandsWithWebClient) .filter(Brand::isNotEmptyBrand) .distinct(); return brandRepo.saveAll(brandStream); }
说明
方案二无需将所有品牌数据加载到内存中,而是以流式的方式逐个(或批量)写入数据库,在数据量较大时性能和内存占用表现更优,是响应式编程的最佳实践。
内容的提问来源于stack exchange,提问作者richardstone
相关产品推荐
相关产品推荐

