Reactor中Flux使用map()不推送元素但flatMap()生效问题咨询
核心原因
map与flatMap对返回值的处理逻辑存在本质差异,结合doWork返回Mono的场景,直接导致了两段代码的表现不同:
map为同步转换操作符,仅将传入函数的返回值作为普通元素向下游传递,不会对返回值做任何订阅、展开操作。使用map(this::doWork)时,流中传递的元素是doWork返回的Mono实例本身,这类Publisher实例从未被订阅,内部逻辑不会触发,自然不会产生任何事件。flatMap为异步流展开操作符,要求传入函数返回Publisher类型(Mono/Flux均为Publisher子类),会自动订阅每个返回的Publisher,将其内部产生的事件展开到上游流中继续传递,因此doWork的逻辑可以正常执行,事件流转符合预期。
两段代码的具体执行逻辑
使用map的异常代码执行流
final Flux<GroupedFlux<String, TData>> groupedFlux = flux.groupBy(Event::getPartitionKey); groupedFlux.subscribe(g -> g.delayElements(Duration.ofMillis(100)) .map(this::doWork) .doOnError(throwable -> log.error("error: ", throwable)) .onErrorResume(e -> Mono.empty()) .subscribe());
- 分组后的每个
GroupedFlux元素经过delayElements延迟后传入map,调用doWork得到未订阅的Mono实例 Mono实例作为普通对象向下游传递,直到流结束- 全程无代码订阅该
Mono,doWork内部逻辑完全不触发,无任何事件输出
使用flatMap的正常代码执行流
final Flux<GroupedFlux<String, TData>> groupedFlux = flux.groupBy(Event::getPartitionKey); groupedFlux.subscribe(g -> g.delayElements(Duration.ofMillis(100)) .flatMap(this::doWork) .doOnError(throwable -> log.error("error: ", throwable)) .onErrorResume(e -> Mono.empty()) .subscribe());
- 分组后的每个元素延迟后传入
flatMap,调用doWork得到Mono实例 flatMap自动订阅该Mono,触发doWork内部逻辑执行Mono产生的事件被flatMap展开到上级流中,后续操作符可以正常接收事件处理,逻辑符合预期
补充注意点
当转换函数返回值为Publisher类型时,必须使用flatMap系列操作符,map仅适用于返回普通POJO的同步转换场景。
内容的提问来源于stack exchange,提问作者ab m
相关产品推荐
相关产品推荐

