为何Flux中delaySequence()结合concatMap()表现如delayElements()?
为何delaySequence()与concatMap()结合时表现得和delayElements()一样,会在每个元素之间产生延迟?而delaySequence()与flatMapSequential()结合时则符合预期。
Flux中有两个用于延迟处理的方法:
delayElements()将
Flux的每个元素延迟指定Duration时长。delaySequence()将
Flux整体在时间上延迟指定Duration时长。与delayElements(Duration)不同,元素会在发射时整体前移,元素间的延迟始终与源一致(仅第一个元素相对于订阅事件有明显延迟)。
使用该操作符时,每秒发射10个元素的源在设置1秒的delaySequence后,仍会以10Hz的频率发射,仅初始有1秒的停顿。而delayElements(Duration)会将发射频率变为1Hz。
当调用.delaySequence(1s).flatMapSequential()时,表现符合上述描述。
但调用.delaySequence(1s).concatMap()时,却表现得像delayElements(),原本10Hz的源最终变成1Hz的发射频率。
我在flatMapSequential()和concatMap()的文档中找不到能解释这种差异的内容。
以下Java代码示例说明了该现象。
delayElements() + flatMapSequential()
@Test void delayElements_flatMapSequential() { log.info("delayElements_flatMapSequential"); Flux.just("one", "two", "three") .delayElements(Duration.ofSeconds(1)) .flatMapSequential(e -> { log.info(e); return Mono.just(e); }) .blockLast(); }
运行结果符合预期,既有初始延迟,也有元素间的延迟:
15:04:03.771 delayElements_flatMapSequential 15:04:04.975 one 15:04:05.976 two 15:04:06.989 three
delaySequence() + flatMapSequential()
@Test void delaySequence_flatMapSequential() { log.info("delaySequence_flatMapSequential"); Flux.just("one", "two", "three") .delaySequence(Duration.ofSeconds(1)) .flatMapSequential(e -> { log.info(e); return Mono.just(e); }) .blockLast(); }
运行结果符合预期,仅有初始延迟,元素间无延迟:
15:04:10.036 delaySequence_flatMapSequential 15:04:11.047 one <- 预期延迟 15:04:11.047 two <- 无延迟 15:04:11.047 three <- 无延迟
delaySequence() + concatMap()
@Test void delaySequence_concatMap() { log.info("delaySequence_concatMap"); Flux.just("one", "two", "three") .delaySequence(Duration.ofSeconds(1)) .concatMap(e -> { log.info(e); return Mono.just(e); }) .blockLast(); }
运行结果不符合预期,既有初始延迟,也有元素间的延迟:
15:04:06.997 delaySequence_concatMap 15:04:08.016 one <- 预期延迟 15:04:09.027 two <- 意外延迟 15:04:10.033 three <- 意外延迟
核心差异源于两个操作符的预取策略和delaySequence的内部实现逻辑:
预取策略的区别
flatMapSequential默认预取32个元素,会一次性向上游请求足够多的元素,delaySequence可以提前获取到源的所有元素,按文档描述的逻辑,在整体延迟1秒后,将所有元素按源的原始间隔(这里Flux.just的元素间隔为0)依次发射,因此元素间无额外延迟。concatMap默认仅预取1个元素,它会先请求1个元素,等当前元素对应的内部流处理完成后,才会向上游请求下一个元素。
delaySequence的适配逻辑
delaySequence的设计是基于源元素的原始发射时间来整体偏移,但当下游每次仅请求1个元素时,它无法提前获取后续元素的时间戳,只能在收到新请求时,为当前请求的元素重新计算延迟——这就导致每个元素都被单独延迟了指定时长,表现得和delayElements一致。
验证方式
你可以修改concatMap的预取参数,将代码改为:
.concatMap(e -> { log.info(e); return Mono.just(e); }, 32) // 设置预取数量为32
此时concatMap会一次性请求足够多的元素,delaySequence的表现就会和flatMapSequential一致,仅初始有1秒延迟,元素间无额外间隔。
内容的提问来源于stack exchange,提问作者Honza Zidek

