【问题标题】:Resilience4j rate limiter is not working properly in project reactor?Resilience4j 速率限制器在项目反应堆中无法正常工作?
【发布时间】:2021-04-27 07:06:07
【问题描述】:

我目前正在研究弹性 4j 库,由于某种原因,以下代码无法按预期工作:

@Test
public void testRateLimiterProjectReactor()
{
    // The configuration below will allow 2 requests per second and a "timeout" of 2 seconds.
    RateLimiterConfig config = RateLimiterConfig.custom()
                                                .limitForPeriod(2)
                                                .limitRefreshPeriod(Duration.ofSeconds(1))
                                                .timeoutDuration(Duration.ofSeconds(2))
                                                .build();

    // Step 2.
    // Create a RateLimiter and use it.
    RateLimiterRegistry registry = RateLimiterRegistry.of(config);
    RateLimiter rateLimiter = registry.rateLimiter("myReactorServiceNameLimiter");

    // Step 3.
    Flux<Integer> flux = Flux.from(Flux.range(0, 10))
                             .transformDeferred(RateLimiterOperator.of(rateLimiter))
                             .log()

        ;

    StepVerifier.create(flux)
                .expectNextCount(10)
                .expectComplete()
                .verify()
    ;
}

根据官方示例herehere,这应该将每秒request() 限制为2 元素。但是,日志显示它会立即获取所有元素:

15:08:24.587 [main] DEBUG reactor.util.Loggers - Using Slf4j logging framework
15:08:24.619 [main] INFO reactor.Flux.Defer.1 - onSubscribe(RateLimiterSubscriber)
15:08:24.624 [main] INFO reactor.Flux.Defer.1 - request(unbounded)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(0)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(1)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(2)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(3)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(4)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(5)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(6)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(7)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(8)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onNext(9)
15:08:24.626 [main] INFO reactor.Flux.Defer.1 - onComplete()

我看不出有什么问题?

【问题讨论】:

  • resilience4j 速率限制器限制订阅数量而不是元素数量
  • Flux 也有内置的速率限制器作为替代:projectreactor.io/docs/core/release/api/reactor/core/publisher/…
  • 感谢您的确认!是的,过去几个小时我一直在查看代码,并在 javadocs 中也注意到了这一点。我对limitRate 的问题是它是固定的——一旦你组装了Flux,你就不能改变费率,这就是我尝试 Resilience4j 的原因。这是stackoverflow.com/questions/67133878/… 的延续,我目前正在研究一旦太多错误/重试开始堆积,您可以通过某种方式逐渐降低limitRate。有什么想法吗?
  • P.S.随意回答最初的问题(它确实限制了订阅)。

标签: project-reactor resilience4j


【解决方案1】:

正如上面 cmets 中已经回答的那样,RateLimiter 跟踪订阅的数量,而不是元素。要实现对元素的速率限制,您可以使用 limitRate(和 buffer + delayElements)。 例如,

        Flux.range(1, 100)
                .delayElements(Duration.ofMillis(100)) // to imitate a publisher that produces elements at a certain rate
                .log()
                .limitRate(10) // used to requests up to 10 elements from the publisher
                .buffer(10) // groups integers by 10 elements
                .delayElements(Duration.ofSeconds(2)) // emits a group of ints every 2 sec
                .subscribe(System.out::println);

【讨论】:

    猜你喜欢
    • 2016-12-04
    • 2020-08-05
    • 2021-04-06
    • 1970-01-01
    • 2015-02-12
    • 2017-11-09
    • 2019-06-17
    • 1970-01-01
    • 2015-08-09
    相关资源
    最近更新 更多