【问题标题】:Flux from iterable how to handle back pressure来自可迭代的通量如何处理背压
【发布时间】: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


【解决方案1】:

如果不能支持许多数据以在处理后发送对新数据的请求,handleAlarm 作业是否有背压 1. 如果你不能有 backPressure to handleAlarm 你可以添加延迟

Flux.fromIterable(alarms)
        .limitRate(parallelism)
        .delayElements(Duration.ofMillis(10))
        .doOnNext(reactiveHandleDataService::handleAlarm)
        .subscribeOn(Schedulars.elastic)

;

如果你想要背压,为什么要并行运行

【讨论】:

  • 我需要一些东西,也许不是背压解决方案,它允许控制数据从通量到下一步的流动。可能我只需要更好地调整并行性,我不知道,我是 reactor 的新手,所以我正在尝试评估所有机会。
  • 在我的情况下,如果句柄警报可以为 10 英里处理 1 个项目,则您设置延迟发射时间,您将在 10 英里内仅发射 1 个项目并且不会丢失任何数据@Stefania
  • 如果警报的大小大于“parallel”,则“limitRate”不允许通量细化警报列表,将其丢弃。
猜你喜欢
  • 1970-01-01
  • 2022-01-04
  • 1970-01-01
  • 2018-12-14
  • 2023-01-20
  • 2023-03-31
  • 1970-01-01
  • 2011-02-18
  • 2018-09-29
相关资源
最近更新 更多