Reactor中GroupedFlux调用delayElements出现跨全组统一延迟问题
问题根因
- 嵌套
subscribe是Reactor的反模式,会导致上下游流量控制失效,且所有GroupedFlux的处理默认共享上游Flux的发射线程,无法实现分组独立调度 delayElements默认使用公共的Schedulers.parallel()调度器,加上上游元素是串行发射到各分组的,会导致所有分组的延迟逻辑排队生效,最终出现全局每100ms仅输出1条消息的现象
解决方案
不要使用嵌套订阅,改用flatMap合并所有分组的处理流,同时为每个分组单独指定调度器,保证不同分组的延迟逻辑互不干扰,修正后代码如下:
final Flux<GroupedFlux<String, TData>> groupedFlux = flux.groupBy(Event::getPartitionKey); groupedFlux // flatMap并发参数根据实际最大分组数调整,需保证大于等于你的分组数量 .flatMap(grouped -> grouped // 每个分组切换到独立调度器,IO密集型场景用boundedElastic,CPU密集型可换用parallel .publishOn(Schedulers.boundedElastic()) .delayElements(Duration.ofMillis(100)) .flatMap(this::doWork) .doOnError(throwable -> log.error("error: ", throwable)) .onErrorResume(e -> Mono.empty()), 32 ) .subscribe();
效果说明
修改后每个分组的delayElements会独立计算延迟,3个分组的场景下每100ms会同时输出3条不同分组的消息,符合预期。
内容的提问来源于stack exchange,提问作者ab m
相关产品推荐
相关产品推荐

