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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 15:55:14