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

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通过两个信号流控制窗口的开启和关闭:

  1. 第一个参数:窗口开启信号流,每当该流发射元素时,打开一个新窗口
  2. 第二个参数:窗口关闭信号生成函数,针对每个开启信号,返回一个流,当该流发射元素时关闭当前窗口

示例代码:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 06:12:11