Flux中delaySequence上游的doAfterTerminate为何先于整条流完成触发?
问题解答
核心原因:delaySequence的实现机制导致
delaySequence的作用是将上游流发出的所有信号(onNext/onComplete/onError)整体缓存,等指定延迟时间结束后,再统一发给下游。它的执行逻辑是:
- 立即订阅上游流,接收上游发出的所有信号
- 把所有信号存入内部缓存,直到上游流终止(收到
onComplete/onError) - 启动定时任务,等待指定延迟时间
- 延迟时间到后,把缓存的所有信号按顺序发给下游
带delaySequence的测试现象解释
你代码中第一个doAfterTerminate是绑定在原始上游Flux.fromArray的终止回调:
- 原始
Flux.fromArray会立刻把1、2、3三个元素和onComplete信号发给下游的delaySequence,自身就终止了,所以第一个doAfterTerminate的回调会立刻执行,打印Finished processing batch! - 此时
delaySequence还在等待1秒的延迟结束,才会把缓存的信号发给下游 - 1秒后
delaySequence把三个元素发往下游,触发doOnNext打印,最后发onComplete触发第二个doAfterTerminate的回调,和你给出的输出顺序完全一致。
删除delaySequence后的测试现象解释
去掉delaySequence后,流没有了缓存环节,信号会直接从上游传递到下游:
- 三个元素依次从上游发出,直接传递到
doOnNext打印Done x - 所有元素发完后,上游发出
onComplete信号,沿着操作符链向下传递 - 这里要注意
doAfterTerminate的回调执行规则:收到终止信号后,先把终止信号传递给下游,再执行自身的回调,所以下游的第二个doAfterTerminate的回调会先执行(打印Finished v2),上游的第一个doAfterTerminate的回调后执行(打印Finished processing batch!),和你给出的输出完全匹配。
调整方案
如果你想让第一个doAfterTerminate在所有元素处理完成后再执行,把它挪到delaySequence的下游即可,它就会绑定到包含延迟逻辑的整个流的终止事件上。
内容的提问来源于stack exchange,提问作者Leon Patmore
相关产品推荐
相关产品推荐

