【发布时间】:2021-04-01 21:09:00
【问题描述】:
我需要对具有高基数的无限 Flux 进行分组。
例如:
- 组键是域 url
- 对一个域的调用应严格按顺序进行(下一个调用发生在前一个调用完成之后)
- 对不同域的调用应该是并发的
- 具有相同键 (url) 的项目之间的时间间隔未知,但预计具有突发性质。几个项目在短时间内发出,然后长时间暂停,直到下一组。
queue
.groupBy(keyMapper, groupPrefetch)
.flatMap(
{ group ->
group.concatMap(
{ task -> makeSlowRemoteCall(task) },
0
)
.takeUntil { remoteCallResult -> remoteCallResult == DONE }
.timeout(groupTimeout, Mono.empty())
.then()
}
, concurrency
)
我在两种情况下取消组:
-
makeSlowRemoteCall()结果表明该组中很有可能在不久的将来不会有新项目。 -
在
groupTimeout期间不会发出下一项。我使用timeout(timeout, fallback)变体来抑制 TimeoutException 并允许 flatMap 的内部发布者成功完成。
我希望将来可能使用相同键的项目来创建新的 GroupedFlux 并使用相同的 flatMap 内部管道进行处理。
但是,如果 GroupedFlux 在我取消时还有剩余的未请求项会发生什么?
groupBy 操作员是否使用相同的密钥将它们重新排队到新组中,否则它们将永远丢失。如果以后什么是解决我的问题的正确方法。我也不确定在这种情况下是否需要将 concatMap() prefetch 设置为 0。
【问题讨论】:
-
如果发生超时,我不清楚你想发生什么——在这种情况下,会抛出一个错误,整个链将因此停止(尚未处理的元素不会在某处重新排队,它们会丢失。)你想在某处记录剩余的元素,还是在超时后对它们执行一些其他操作,或者完全是其他什么?
-
我编辑了我的问题以使其更清楚。请注意使用
timeout(timeout, fallback)变体来抑制超时错误。它应该取消源通量组,但允许 flatMap 内部发布者成功完成。基本上我希望剩余的项目由相同的 concatMap 管道处理,但作为新的 GroupedFlux。我使用超时的原因是从 flatMap 内部清除空闲组。 -
啊,我明白了。我认为您很不走运-一旦发生超时,整个
GroupedFlux链就会被空的Mono替换,您无法恢复实际的组。话虽如此,听起来您正在使用timeout()作为清除组的优化,这有点奇怪——除非组的数量变得非常大,否则您不需要这样做。如果是这种情况,您可能最好通过为每个域使用一个新的隔板实例来使用弹性 4j (resilience4j.readme.io/docs/bulkhead) 中的隔板模式之类的东西。 -
(续) ...如果我理解正确的话,以这种方式使用隔板应该可以让您完全避免分组 - 您可以在标准
flatMap中进行慢速远程调用,但是使用舱壁名称作为键。 -
以下 javadoc 让我担心“请注意,当 <...> 标准产生大量组时,groupBy 最适合低基数的组,如果组不是,则可能导致挂起下游适当消耗”。我将拥有多达 10k 个可能的键,每秒有 100-1000 个入站元素。每组最多空闲 30 秒,直到终端元素到达,我可以取消它。不挂 flatMap 听起来够好吗?
标签: project-reactor