【发布时间】:2018-12-21 11:38:07
【问题描述】:
我需要处理我的通量中的背压,它接收一个对象列表作为输入。列表的大小从几百到几十万个元素不等。 实际代码是:
Flux.fromIterable(alarms)
.limitRate(parallelism)
.parallel(parallelism)
.runOn(Schedulers.elastic(), bufferSize)
.doOnNext(reactiveHandleDataService::handleAlarm)
;
参数“limitRange”只是强制拒绝超过一定大小的列表,这是我不想要的。我需要将收到的所有数据提供给 reactiveHandleDataService,我不能丢失消息。
在这种情况下,我该如何处理背压?我没有找到太多解释问题的例子,尤其是使用可迭代作为源代码。
我使用 Californium-SR3 作为反应器的发布版本,这是 Spring Boot 应用程序的一部分。
【问题讨论】:
-
该项目列表来自哪里?如果您已经拥有
List中的所有内容,那么应用backpressure 有什么意义? -
我需要背压以不使方法 handleAlarm 过载,从而不使 handleAlarm 中使用的连接池饱和。
标签: java spring-boot reactor