You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Resilience4j与Reactor RetryWhen协同异常:无法按10RPS分发100请求

问题描述

调用远程服务时希望限制在10 RPS,因此配置了Resilience4j Rate Limiter,并添加retryWhen处理RequestNotPermitted错误,打算在令牌允许时重试。但目前重试未实际多次执行,无法获取ID大于10的请求结果。

示例代码

public static void main(final String[] args) {

    final RateLimiterConfig config = RateLimiterConfig.custom()
        .limitForPeriod(10)
        .limitRefreshPeriod(Duration.ofSeconds(1))
        .timeoutDuration(Duration.ofMillis(0))
        .build();

    final RateLimiterRegistry rateLimiterRegistry = RateLimiterRegistry.of(config);
    final RateLimiter rateLimiter = rateLimiterRegistry.rateLimiter("exampleService");

    final Flux<Integer> ids = Flux.range(1, 100); // 生成1到100的ID流

    ids.flatMapSequential(id -> Mono.defer(() -> invokeRemoteService(id)
        .transformDeferred(RateLimiterOperator.of(rateLimiter))
        .retryWhen(retrySpec(id))))
        .doOnNext(response -> log.info("Response: {}", response))
        .subscribe();
}

private static Mono<String> invokeRemoteService(final Integer integer) {
    return Mono.just(integer + "");
}

private static RetryBackoffSpec retrySpec(final Integer id) {
    return Retry.fixedDelay(10, Duration.ofMillis(1000))
        .filter(throwable -> {
            final boolean retry = throwable instanceof RequestNotPermitted;
            log.info("Retrying {}? {}", id, retry);
            return retry;
        });
}

输出结果

11:07:04.100 [main] INFO com.sample.Resilience -- Response: 1
11:07:04.105 [main] INFO com.sample.Resilience -- Response: 2
11:07:04.106 [main] INFO com.sample.Resilience -- Response: 3
11:07:04.107 [main] INFO com.sample.Resilience -- Response: 4
11:07:04.107 [main] INFO com.sample.Resilience -- Response: 5
11:07:04.107 [main] INFO com.sample.Resilience -- Response: 6
11:07:04.108 [main] INFO com.sample.Resilience -- Response: 7
11:07:04.108 [main] INFO com.sample.Resilience -- Response: 8
11:07:04.109 [main] INFO com.sample.Resilience -- Response: 9
11:07:04.109 [main] INFO com.sample.Resilience -- Response: 10
11:07:04.117 [main] INFO com.sample.Resilience -- Retrying 11? true
11:07:04.153 [main] INFO com.sample.Resilience -- Retrying 12? true
11:07:04.154 [main] INFO com.sample.Resilience -- Retrying 13? true
...
11:07:04.220 [main] INFO com.sample.Resilience -- Retrying 99? true
11:07:04.221 [main] INFO com.sample.Resilience -- Retrying 100? true
问题根源与修复方案

问题原因

  1. 超时时间配置错误:timeoutDuration设为0毫秒,导致令牌耗尽后RateLimiter直接拒绝请求,没有等待令牌刷新的窗口,所有后续请求瞬间触发RequestNotPermitted。
  2. 重试逻辑位置错误:重试操作嵌套在Mono.defer内部,仅对远程调用生效,而非包含限流的完整逻辑。
  3. 主线程提前终止:使用subscribe()是非阻塞调用,主线程会快速遍历所有ID并触发重试日志,随后直接结束,不会等待重试延迟后的实际请求执行。

修复步骤

1. 调整RateLimiter超时配置

设置合理的超时时间,让RateLimiter在令牌不足时等待刷新:

final RateLimiterConfig config = RateLimiterConfig.custom()
    .limitForPeriod(10)
    .limitRefreshPeriod(Duration.ofSeconds(1))
    .timeoutDuration(Duration.ofSeconds(1)) // 修改为1秒超时
    .build();

2. 调整重试逻辑位置

将retryWhen移到flatMapSequential的Mono外层,确保限流+远程调用的完整逻辑被重试:

ids.flatMapSequential(id -> Mono.defer(() -> invokeRemoteService(id)
        .transformDeferred(RateLimiterOperator.of(rateLimiter)))
    .retryWhen(retrySpec(id))) // 重试包含限流的整个流程

3. 阻塞主线程等待执行完成

替换subscribe()为blockLast(),让主线程等待所有请求处理完毕:

.doOnNext(response -> log.info("Response: {}", response))
.blockLast();

完整修复代码

public static void main(final String[] args) {

    final RateLimiterConfig config = RateLimiterConfig.custom()
        .limitForPeriod(10)
        .limitRefreshPeriod(Duration.ofSeconds(1))
        .timeoutDuration(Duration.ofSeconds(1))
        .build();

    final RateLimiterRegistry rateLimiterRegistry = RateLimiterRegistry.of(config);
    final RateLimiter rateLimiter = rateLimiterRegistry.rateLimiter("exampleService");

    final Flux<Integer> ids = Flux.range(1, 100);

    ids.flatMapSequential(id -> Mono.defer(() -> invokeRemoteService(id)
            .transformDeferred(RateLimiterOperator.of(rateLimiter)))
        .retryWhen(retrySpec(id)))
        .doOnNext(response -> log.info("Response: {}", response))
        .blockLast();
}

private static Mono<String> invokeRemoteService(final Integer integer) {
    return Mono.just(integer + "");
}

private static RetryBackoffSpec retrySpec(final Integer id) {
    return Retry.fixedDelay(10, Duration.ofMillis(1000))
        .filter(throwable -> {
            final boolean retry = throwable instanceof RequestNotPermitted;
            log.info("Retrying {}? {}", id, retry);
            return retry;
        });
}

修复效果

配置调整后,RateLimiter会在令牌不足时等待刷新,重试逻辑会在限流后延迟1秒重试,主线程会一直等待所有100个请求完成,最终能看到ID从1到100的响应依次输出,符合10 RPS的限制。

内容的提问来源于stack exchange,提问作者codependent

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 01:43:10