【问题标题】:Flux.range waits to emit more element once 256 elements are reached一旦达到 256 个元素,Flux.range 就会等待发射更多元素
【发布时间】:2020-12-01 10:23:52
【问题描述】:

我写了这段代码:

Flux.range(0, 300)
            .doOnNext(i -> System.out.println("i = " + i))
            .flatMap(i -> Mono.just(i)
                            .subscribeOn(Schedulers.elastic())
                            .delayElement(Duration.ofMillis(1000))
            )
            .doOnNext(i -> System.out.println("end " + i))
            .blockLast();

运行它时,第一个System.out.println 表明 Flux 在第 256 个元素处停止发射数字,然后它等待旧的完成后再发射新的。

为什么会这样?
为什么是 256?

【问题讨论】:

    标签: java project-reactor flux


    【解决方案1】:

    为什么会这样?

    flatMap 运算符可以描述为以下运算符(改写自 javadoc):

    1. 订阅其内部热切
    2. 不保留元素的顺序。
    3. 让来自不同内部的值交错。

    对于这个问题,第一点很重要。项目反应堆限制 通过concurrency 参数实现的inner 序列的数量。

    虽然flatMap(mapper) 使用默认参数,但flatMap(mapper, concurrency) 重载显式接受此参数。

    flatMaps javadoc 将参数描述为:

    concurrency 参数允许控制可以并行订阅和合并多少个 Publisher

    考虑以下代码,使用concurrency = 500

    Flux.range(0, 300)
            .doOnNext(i -> System.out.println("i = " + i))
            .flatMap(i -> Mono.just(i)
                            .subscribeOn(Schedulers.elastic())
                            .delayElement(Duration.ofMillis(1000)),
                    500
    //         ^^^^^^^^^^
            )
            .doOnNext(i -> System.out.println("end " + i))
            .blockLast();
    

    在这种情况下没有等待:

    i = 297
    i = 298
    i = 299
    end 0
    end 1
    end 2
    

    相比之下,如果您将 1 传递为 concurrency,输出将类似于:

    i = 0
    end 0
    i = 1
    end 1
    

    在发出下一个元素之前等待一秒钟。

    为什么是 256?

    256 是flatMap 并发的默认值。

    看看Queues.SMALL_BUFFER_SIZE

    public static final int SMALL_BUFFER_SIZE = Math.max(16,
            Integer.parseInt(System.getProperty("reactor.bufferSize.small", "256")));
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-05-31
      • 1970-01-01
      • 2022-11-05
      • 2021-11-28
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多