【问题标题】:Reactor 3.x - limit the time of a groupBy FluxReactor 3.x - 限制 groupBy Flux 的时间
【发布时间】:2018-06-02 06:40:40
【问题描述】:

是否有任何方法可以强制由 groupBy() 生成的 Flux 在一段时间后完成(或者类似地,限制“打开”组的最大数量)而不管上游的完整性?我有类似以下内容:

Flux<Foo> someFastPublisher;

someFastPublisher
  .groupBy(f -> f.getKey())
  .delayElements(Duration.ofSeconds(1)) // rate limit each group
  .flatMap(g -> g) // unwind the group
  .subscribe()
;

我遇到了 Flux 挂起的情况,假设是因为组数大于 flatMap 的并发性。我可以增加flatMap 的并发性,但是没有简单的方法来判断最大可能大小是多少。相反,我知道由Foo.key 分组的Foo 将在时间/发布顺序上彼此接近,并且宁愿在 groupBy Flux 与 flatMap 并发上使用某种时间窗口(并最终有两个不同的组和相同的key() 没什么大不了的)。

我猜groupBy Flux 在someFastPubisher onCompletes 之前不会 onComplete - 即 Flux 已移交给 flatMap 只是保持“开放”(尽管他们不太可能获得新事件) .

我可以通过在 groupBy 中预取 Integer.MAXInteger.MAXing 来解决这个问题 - 但是有没有办法控制组的“生命”?

【问题讨论】:

    标签: project-reactor


    【解决方案1】:

    是的:您可以将take(Duration) 应用到群组,以确保它们提前关闭,之后将打开具有相同密钥的新群组:

    source.groupBy(v -> v.intValue() % 2)
          .flatMap(group -> group
                  .take(Duration.ofMillis(1000))
                  .count()
                  .map(c -> "group " + group.key() + " size = " + c)
          )
          .log()
          .blockLast();
    

    【讨论】:

    • 考虑到谁回答了它,我会相信这个作品 :) 但我很难想象发生了什么。 take() 是否会导致 groupedBy Flux 在持续时间后 onComplete ?和/或如何向 GroupedFlux 上游发出信号?
    • 没关系 - 我想我明白了。不是上游 groupBy 发出信号 - 而是 flatMap 看到 take() Flux(副本)而不是 GroupedFlux Flux
    • 我宁愿使用groupByDuration 参数作为语法糖的方法。没有提示就很难采用这种方法。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-09-10
    • 2021-11-23
    • 2019-09-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多