【问题标题】:Not able to exhibit backpressure with spring web reactive无法通过弹簧网反应式表现出背压
【发布时间】:2017-07-04 23:57:34
【问题描述】:

我正在尝试使用 spring-web-reactive 来展示背压,就像这里使用 akka 显示的方式一样 - https://www.youtube.com/watch?v=oS9w3VenDW0 (在 28:20 到 29:20 之间观看)。

为了试用,我使用了来自 github https://github.com/bclozel/spring-boot-web-reactive 的以下示例项目

在设置项目后,我在 HomeController.java 中添加了一个新端点,如下所示:

@RequestMapping(value = "/longflux",produces = "application/stream+json")
public Flux<Long> longFlux(){
    return Flux.interval(Duration.ofMillis(10)).log();
}

现在,如果我尝试卷曲这个端点,然后使用 (CTRL+z) 将其挂起,那么一旦 tcp 缓冲区被填满并且服务器应该停止发出事件,背压应该会启动。

但是,在某个时间后暂停 curl 命令会引发以下异常:

2017-02-16 08:49:48.480 ERROR 3500 --- [        timer-1] reactor.Flux.Interval.4                  : onError(reactor.core.Exceptions$OverflowException: Could not emit value 2578 due to lack of requests)
2017-02-16 08:49:48.481 ERROR 3500 --- [        timer-1] reactor.Flux.Interval.4                  : 
reactor.core.Exceptions$OverflowException: Could not emit value 2578 due to lack of requests
    at reactor.core.Exceptions.failWithOverflow(Exceptions.java:151) ~[reactor-core-3.0.4.RELEASE.jar:3.0.4.RELEASE]
    at reactor.core.publisher.FluxInterval$IntervalRunnable.run(FluxInterval.java:98) ~[reactor-core-3.0.4.RELEASE.jar:3.0.4.RELEASE]
    at reactor.core.scheduler.SingleTimedScheduler$TimedPeriodicScheduledRunnable.run(SingleTimedScheduler.java:394) ~[reactor-core-3.0.4.RELEASE.jar:3.0.4.RELEASE]
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[na:1.8.0_121]

我无法理解为什么请求在 curl 命令暂停后的某个时间异常终止(在 spring-web-reactive 实现中),而在 akka 示例中(如 youtube 链接所示)服务器停止发布tcp 缓冲区满后的事件。

【问题讨论】:

    标签: java spring project-reactor spring-webflux


    【解决方案1】:

    Flux.interval 是一个特例,因为它是一个热源,而且 Reactor 不会缓冲时间;这意味着如果您的请求周期由于背压而变慢并且您的间隔源产生速度更快,Reactor 将发出错误信号。

    您可以使用 .onBackpressureDrop() 运算符更新此示例,以在出现背压的情况下降低间隔。这应该符合预期。

    有很多方法可以说明背压,包括:

    • 使用delay 运算符延迟订阅
    • 模拟多个慢速客户端(带宽和延迟)

    【讨论】:

    • 感谢布赖恩的指点。我现在能够使用下面的源代码来说明背压。这次我没有使用热源,而是创建了一个列表,然后使用流创建源。 @RequestMapping(value = "/longflux",produces = "application/stream+json") public Flux longFlux(){ List list = new ArrayList(); for(int i=0;i
    猜你喜欢
    • 2016-09-20
    • 2018-07-03
    • 2019-04-30
    • 1970-01-01
    • 1970-01-01
    • 2020-09-19
    • 2021-08-12
    • 2023-01-27
    • 1970-01-01
    相关资源
    最近更新 更多