引言响应式编程(Reactive Programming)在构建高并发、异步数据流应用时展现出强大的优势。Reactor作为 Spring WebFlux 的底层实现,提供了Flux和Mono两种核心类型,分别代表 0…N 个和 0…1 个元素的异步序列。本文将带领读者从六个不同来源的数据流出发,完成一个完整的实战任务:分组计算奇数流的逐项乘积和偶数流的逐项和,并解决实际编码中遇到的典型问题。一、有限流与无限流在 Reactor 中,Flux可以是有限的,也可以是无限的。有限流:在发出固定数量的元素后,会触发onComplete信号结束流。例如Flux.range(1, 10)产生 1~10 共 10 个元素,随后完成。无限流:理论上永远不会完成,除非被取消或发生错误。例如Flux.generate可以无限生成元素,通常需要配合.take(n)截取前 n 项来使用。理解二者的区别对于后续的流操作至关重要,因为无限流若不加限制,会导致下游无限等待,而有限流则会自然结束。二、六个核心数据流定义我们将创建六个流,分别标记为flux1到flux6,其中奇数流为flux1、flux3、flux5,偶数流为flux2、flux4、flux6。importreactor.core.publisher.Flux;importreactor.core.publisher.Flux;importjava.math.BigInteger;importjava.time.Duration;importjava.util.List;publicclassFluxDemo{// ---------- 奇数流 ----------// flux1:三角数(1, 3, 6, 10, 15, ...)——无限流staticFluxLongflux1=Flux.generate(()-newlong[]{1,0},// 状态:[当前数, 累计和](state,sink)-{longcurrent=state[0];longsum=state[1];sum+=current;sink.next(sum);state[0]=current+1;state[1]=sum;returnstate;});// flux3:有限列表 [-2, 4, 2, -4, 8] ——有限流,共5项staticFluxIntegerflux3=Flux.fromIterable(List.of(-2,4,2,-4,8)).delayElements(Duration.ofMillis(10)).doFinally(signalType-System.out.println("flux3输出完成: "+signalType));// flux5:20以内的素数(2, 3, 5, 7, 11, 13, 17, 19)——有限流,共8项staticFluxIntegerflux5=Flux.range(2,18).filter(n-{for(inti=2;i*i=n;i++){if(n%i==0)returnfalse;}returntrue;});// ---------- 偶数流 ----------// flux2:1~10的整数 ——有限流,共10项staticFluxIntegerflux2=Flux.range(1,10).delayElements(Duration.ofMillis(10)).doFinally(signalType-System.out.println("flux2输出完成: "+signalType));// flux4:阶乘(1!, 2!, 3!, ...)——无限流staticFluxBigIntegerflux4=Flux.generate(()-newBigInteger[]{BigInteger.ONE,BigInteger.ONE},(state,sink)-{BigIntegern=state[0];BigIntegerfact=state[1];sink.next(fact);BigIntegernextN=n.add(BigInteger.ONE);BigIntegernextFact=fact.multiply(nextN);returnnewBigInteger[]{nextN,nextFact};});// flux6:2的幂(1, 2, 4, 8, ...)——无限流staticFluxIntegerflux6=Flux.generate(()-1,(state,sink)-{sink.next(state);returnstate*2;});}流解析flux1(三角数):使用Flux.generate维护状态数组[current, sum],每次累加current并递增current,得到无限三角数序列。flux3:从固定列表生成,添加了delayElements模拟延迟,并附带doFinally便于观察结束状态。flux5(素数):Flux.range(2,18)产生 2~19,filter中嵌套素数判断,仅保留素数。flux2:Flux.range(1,10)产生 1~10,同样有延迟和doFinally。flux4(阶乘):使用BigInteger避免溢出,状态为[n, fact],每次发出fact并更新为(n+1)!。flux6(2的幂):最简单的generate,状态即当前值,每次发出后乘以2。三、分组与公共长度计算奇数流组:flux1(无限)、flux3(5项)、flux5(8项) → 最少项数为5,记作n = 5。偶数流组:flux2(10项)、flux4(无限)、flux6(无限) → 最少项数为10,记作m = 10。思考:若流数量不确定,可使用.count().block()获取有限流长度,再取最小值。但本例中长度明确,故直接硬编码。四、打印与聚合逻辑我们需要:分别打印奇数流的前 n 项,并在每个数字前标注序号和流名称。计算奇数流的逐项乘积:即第 i 项 =flux1[i] * flux3[i] * flux5[i],输出结果。同样处理偶数流:先分别打印前 m 项,再计算逐项和 =flux2[i] + flux4[i] + flux6[i]。顺序:先奇数流全部输出,再偶数流全部输出。工具函数:带序号打印privatestaticTvoidprintStreamWithIndex(StringstreamName