Project Reactor中Flux事件顺序疑问:onComplete提前触发
Project Reactor中buffer操作与onComplete触发顺序解析
问题场景
作为Project Reactor新手,在使用buffer(2)操作符时遇到了事件顺序的困惑:
- 示例1:当Flux包含2个元素时,输出符合预期,
onComplete在消费完最后一批元素后触发; - 示例2:当Flux包含3个元素时,
onComplete竟然在消费者处理最后一个结果[2]之前就触发了,这与“onComplete是Flux成功完成时触发”的认知不符。
示例1代码
Flux.fromStream(IntStream.range(0, 2).boxed()) // Two element in the stream .doOnSubscribe(subscription -> out.println("OnSubscribe")) .doOnRequest(l -> out.println("OnRequest")) .doOnNext(subscription -> out.println("onNext")) .doOnEach(subscription -> out.println("onEach")) .doOnComplete(() -> out.println("OnComplete")) .doOnTerminate(() -> out.println("onTerminate")) .doAfterTerminate(() -> out.println("doAfterTerminate")) .doFinally(signalType -> out.println("doFinally")) .doOnCancel(() -> out.println("onCancel")) .doOnError(throwable -> out.println("onError")) .buffer(2) // Buffer size 2 .subscribe(integers -> out.println("Numbers " + integers));
示例1输出
OnSubscribe OnRequest onNext onEach onNext onEach Numbers [0, 1] onEach OnComplete onTerminate doFinally doAfterTerminate
示例2代码
Flux.fromStream(IntStream.range(0, 3).boxed()) // Three elements here !!!!! .doOnSubscribe(subscription -> out.println("OnSubscribe")) .doOnRequest(l -> out.println("OnRequest")) .doOnNext(subscription -> out.println("onNext")) .doOnEach(subscription -> out.println("onEach")) .doOnComplete(() -> out.println("OnComplete")) .doOnTerminate(() -> out.println("onTerminate")) .doAfterTerminate(() -> out.println("doAfterTerminate")) .doFinally(signalType -> out.println("doFinally")) .doOnCancel(() -> out.println("onCancel")) .doOnError(throwable -> out.println("onError")) .buffer(2) // Buffer size 2 .subscribe(integers -> out.println("Numbers " + integers));
示例2输出
OnSubscribe OnRequest onNext onEach onNext onEach Numbers [0, 1] onNext onEach onEach OnComplete onTerminate Numbers [2] doFinally doAfterTerminate
核心原因解析
你看到的差异本质是回调绑定的Flux阶段不同:
- 你注册的
doOnComplete、doOnTerminate等回调,都是绑定在buffer(2)操作符之前的上游Flux(也就是基于Stream创建的原始Flux)上的,而非buffer之后的下游Flux。 buffer操作符的工作逻辑是:- 攒够指定数量的元素后,立即向下游发送一批;
- 当上游Flux触发
onComplete信号时,会把当前缓存中剩余的不足指定数量的元素,打包成最后一批发送给下游。
对应示例2的时间线:
- 上游Flux发送元素0、1 →
buffer攒够2个,发送[0,1]给下游,下游打印Numbers [0,1]; - 上游Flux发送元素2 → 触发
onNext和onEach; - 上游Flux所有元素发送完毕,触发自身的
onComplete信号 → 你注册的doOnComplete、doOnTerminate依次执行; buffer操作符收到上游的onComplete后,把剩余的元素2打包成[2]发送给下游 → 下游打印Numbers [2];- 整个流终止,执行
doFinally和doAfterTerminate。
而示例1中,上游Flux的2个元素刚好凑满buffer(2)的批次,buffer发送[0,1]给下游后,上游才触发onComplete,所以看起来onComplete在消费完最后一批之后执行,这只是巧合的对齐情况。
内容的提问来源于stack exchange,提问作者Guillaume Macke
相关产品推荐
相关产品推荐

