Reactive Java:如何在运行时动态调整Flux窗口大小?
动态调整Reactor Flux窗口大小及重建续处理方案
一、实现动态窗口大小的正确方式
问题分析
你原代码的问题在于window(calcWindowSize())只会在初始化时调用一次calcWindowSize(),窗口大小固定,无法实现动态调整。要实现每次窗口都使用新计算的大小,需要用基于条件或动态信号的窗口操作符。
方案1:使用windowUntil实现动态元素计数窗口
windowUntil可以根据元素级别的条件判断是否关闭当前窗口,适合按动态数量分组的场景。我们可以维护一个计数器,每次窗口结束后重新计算下一个窗口的大小:
import reactor.core.publisher.Flux; import reactor.util.function.MutableInt; import java.time.Duration; import java.util.concurrent.ThreadLocalRandom; public class DynamicWindowExample { public static void main(String[] args) throws InterruptedException { Duration interval = Duration.ofSeconds(1); // 跟踪当前窗口剩余需要收集的元素数量 MutableInt remainingItems = MutableInt.of(0); Flux<String> source = Flux.just("A", "B", "C", "D", "E", "F", "G", "H", "I", "J") .repeat() .delayElements(interval); // 每秒发射一个元素 source .windowUntil(item -> { // 剩余数量为0时,计算新窗口大小 if (remainingItems.get() == 0) { remainingItems.set(calcWindowSize()); } remainingItems.decrement(); // 剩余数量为0时关闭当前窗口 return remainingItems.get() == 0; }) .concatMap(Flux::collectList) .doOnNext(batch -> { System.out.printf("Processed batch: %s at %d%n", batch, System.currentTimeMillis()); }) .subscribe(); Thread.sleep(10_000); } private static int calcWindowSize() { // 实际场景中替换为你的动态计算逻辑 return ThreadLocalRandom.current().nextInt(1, 10); } }
方案2:使用windowWhen实现动态窗口
windowWhen通过两个信号流控制窗口的开启和关闭:
- 第一个参数:窗口开启信号流,每当该流发射元素时,打开一个新窗口
- 第二个参数:窗口关闭信号生成函数,针对每个开启信号,返回一个流,当该流发射元素时关闭当前窗口
示例代码:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.Duration; import java.util.concurrent.ThreadLocalRandom; public class WindowWhenExample { public static void main(String[] args) throws InterruptedException { Duration interval = Duration.ofSeconds(1); Flux<String> source = Flux.just("A", "B", "C", "D", "E", "F", "G", "H", "I", "J") .repeat() .delayElements(interval) .share(); // 共享源,避免多订阅重复发射元素 source .windowWhen( // 窗口开启信号:无限流,初始发射一个信号,之后每次窗口关闭后重复发射 () -> Flux.never().startWith("init").repeat(), // 窗口关闭信号:收集指定数量元素后关闭窗口 openSignal -> source.take(calcWindowSize()).then() ) .concatMap(Flux::collectList) .doOnNext(batch -> { System.out.printf("Processed batch: %s at %d%n", batch, System.currentTimeMillis()); }) .subscribe(); Thread.sleep(10_000); } private static int calcWindowSize() { return ThreadLocalRandom.current().nextInt(1, 10); } }
二、重建Flux并从之前位置续处理
如果因某些限制无法动态调整窗口,可以通过状态保存的方式重建Flux并续处理,核心是记录当前处理的位置,重建时从该位置继续。
示例:基于固定序列的续处理
假设你的元素是固定重复的序列,用原子类保存当前处理的索引,重建时跳过已处理的元素:
import reactor.core.publisher.Flux; import java.time.Duration; import java.util.List; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.atomic.AtomicInteger; public class ReconstructFluxExample { // 保存当前处理的元素索引,线程安全 private static final AtomicInteger currentIndex = new AtomicInteger(0); private static final List<String> ELEMENTS = List.of("A", "B", "C", "D", "E", "F", "G", "H", "I", "J"); public static void main(String[] args) throws InterruptedException { Duration interval = Duration.ofSeconds(1); // 第一次订阅处理 processFlux(interval).subscribe(); Thread.sleep(5_000); // 模拟重建Flux,从之前的位置继续处理 System.out.println("=== 重建Flux,续处理 ==="); processFlux(interval).subscribe(); Thread.sleep(5_000); } private static Flux<List<String>> processFlux(Duration interval) { return Flux.defer(() -> { int startIdx = currentIndex.get(); // 从当前索引开始生成元素流 return Flux.generate( () -> startIdx, (state, sink) -> { int elemIdx = state % ELEMENTS.size(); sink.next(ELEMENTS.get(elemIdx)); currentIndex.incrementAndGet(); return state + 1; } ); }) .delayElements(interval) .window(calcWindowSize()) .concatMap(Flux::collectList) .doOnNext(batch -> { System.out.printf("Processed batch: %s at %d%n", batch, System.currentTimeMillis()); }); } private static int calcWindowSize() { return ThreadLocalRandom.current().nextInt(1, 10); } }
复杂场景扩展
如果元素不是固定序列,可以结合外部存储(如Redis、数据库)保存处理状态(比如偏移量、最后处理的元素标识),重建Flux时从存储中读取状态,跳过已处理的元素。Reactor也提供了checkpoint() API用于故障恢复,但需要配合外部状态存储使用。
内容的提问来源于stack exchange,提问作者krinklesaurus
相关产品推荐
相关产品推荐

