Flux分组并行处理在唯一ID为256时无限挂起问题求助
问题分析与解决
核心原因
你遇到的挂起问题,本质是Reactor groupBy操作的默认分组缓存阈值导致的死锁:
groupBy默认最多同时缓存256个活跃分组,当达到这个数量时,上游会停止发送新元素,直到有分组被消费完毕并销毁。- 测试代码里用了
repeat(),每个分组会持续收到新元素,没有任何分组会进入完成状态,上游被永久阻塞,整个流陷入死锁。 - 当唯一ID数小于256时,分组数没触达阈值,上游正常发送元素;大于256时,
groupBy的缓存会自动扩容,反而不会触发死锁,只有刚好等于256时触发了阈值阻塞逻辑。
解决方案
方案1:调整groupBy的缓存阈值
显式设置groupBy的bufferSize参数,把阈值调大,避免触发上游阻塞:
@ParameterizedTest @ValueSource(ints = {250, 251, 252, 253, 254, 255, 256}) void freezeTest(int uniqueStringsCount) { var scheduler = Schedulers .newBoundedElastic( 1000, 1000, "really-big-scheduler" ); Flux.range(0, uniqueStringsCount) .map(Object::toString) .repeat() .take(50_000) // 调整bufferSize到大于预期的分组数 .groupBy(x -> x, 1024) .parallel() .flatMap(group -> group.concatMap(e -> Mono.fromRunnable(() -> { try { Thread.sleep(0); } catch (InterruptedException ex) { throw new RuntimeException(ex); } }) ) ) .runOn(Schedulers.parallel()) .then() .block(); }
方案2:替换parallel()为带并发控制的flatMap
直接用flatMap控制组间并行度,同时保留组内的concatMap串行处理,这种方式更直观,也能避开groupBy的缓存限制:
@ParameterizedTest @ValueSource(ints = {250, 251, 252, 253, 254, 255, 256}) void freezeTest(int uniqueStringsCount) { var scheduler = Schedulers .newBoundedElastic( 1000, 1000, "really-big-scheduler" ); Flux.range(0, uniqueStringsCount) .map(Object::toString) .repeat() .take(50_000) .groupBy(x -> x) // 用flatMap控制并行处理的分组数量,比如设为1000 .flatMap(group -> group.concatMap(e -> Mono.fromRunnable(() -> { try { Thread.sleep(0); } catch (InterruptedException ex) { throw new RuntimeException(ex); } }) ), 1000 ) .runOn(Schedulers.parallel()) .then() .block(); }
额外说明
- 如果实际业务中分组会有完成的情况(比如某个ID不再产生事件),
groupBy的默认阈值不会有问题;但如果是持续产生事件的分组,必须调整缓存阈值或者换用其他并行控制方式。 Schedulers.parallel()的默认线程数是CPU核心数,若需要更高并行度,可换成Schedulers.newBoundedElastic()或自定义调度器,但核心问题还是groupBy的缓存阈值。
内容的提问来源于stack exchange,提问作者Dmytro Marchuk
相关产品推荐
相关产品推荐

