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
问题根源与修复方案
问题原因
- 超时时间配置错误:
timeoutDuration设为0毫秒,导致令牌耗尽后RateLimiter直接拒绝请求,没有等待令牌刷新的窗口,所有后续请求瞬间触发RequestNotPermitted。 - 重试逻辑位置错误:重试操作嵌套在
Mono.defer内部,仅对远程调用生效,而非包含限流的完整逻辑。 - 主线程提前终止:使用
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
相关产品推荐
相关产品推荐

