【发布时间】: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.MAX 或 Integer.MAXing 来解决这个问题 - 但是有没有办法控制组的“生命”?
【问题讨论】:
标签: project-reactor