You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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操作符的工作逻辑是:
    1. 攒够指定数量的元素后,立即向下游发送一批;
    2. 当上游Flux触发onComplete信号时,会把当前缓存中剩余的不足指定数量的元素,打包成最后一批发送给下游。

对应示例2的时间线:

  1. 上游Flux发送元素0、1 → buffer攒够2个,发送[0,1]给下游,下游打印Numbers [0,1];
  2. 上游Flux发送元素2 → 触发onNext和onEach;
  3. 上游Flux所有元素发送完毕,触发自身的onComplete信号 → 你注册的doOnComplete、doOnTerminate依次执行;
  4. buffer操作符收到上游的onComplete后,把剩余的元素2打包成[2]发送给下游 → 下游打印Numbers [2];
  5. 整个流终止,执行doFinally和doAfterTerminate。

而示例1中,上游Flux的2个元素刚好凑满buffer(2)的批次,buffer发送[0,1]给下游后,上游才触发onComplete,所以看起来onComplete在消费完最后一批之后执行,这只是巧合的对齐情况。


内容的提问来源于stack exchange,提问作者Guillaume Macke

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 01:55:21