Resilience4j RateLimiter在Spring WebFlux中内存耗尽问题求助
问题分析与解决方案
核心问题根源
你当前的内存耗尽问题,本质是上游数据生产速度远大于下游第三方服务的消费速度,导致大量未处理的record对象在内存中堆积。具体原因:
parallel(16)+ 默认并发数的flatMap会让上游R2DBC快速拉取大量记录(80k),即使下游的callCarsService有限流,这些记录也会被缓存等待调用,持续占用内存。- 每个请求单独加限流的方式,无法控制上游整体的数据生产节奏,只能限制单个请求的发送频率,无法从根源解决数据堆积问题。
针对性解决方案
1. 限制flatMap的并发数
flatMap默认并发数是256,远大于你180次/秒的限流值,会导致大量请求排队等待。将并发数调整为接近限流值,减少内存中等待的record数量:
Flux.defer(() -> R2dbcRepository.selectRecords()) //~80k records // ... 其他中间操作 .parallel(16) .runOn(Schedulers.parallel()) // ... 其他中间操作 .flatMap(record -> callCarsService(record), 200) // 并发数设为略大于限流值 .sequential();
2. 全局限流替代单请求限流
将限流逻辑从callCarsService移到主Flux流中,从整体上控制数据处理的速率,同时结合R2DBC的背压机制,避免上游过度拉取数据:
// 定义全局限流器 RateLimiter rateLimiter = RateLimiter.of("global-limit", RateLimiterConfig.custom() .limitRefreshPeriod(Duration.ofSeconds(1)) .limitForPeriod(180) .build()); // 调整主流程 Flux.defer(() -> R2dbcRepository.selectRecords()) // 限制从数据库拉取的速率,避免一次性加载过多记录到内存 .limitRate(1000) .parallel(16) .runOn(Schedulers.parallel()) // 在并行处理后添加全局限流,控制整体处理速率 .transform(RateLimiterOperator.of(rateLimiter)) .flatMap(record -> callCarsService(record), 180) .sequential(); // 修改callCarsService,移除内部限流逻辑 Mono<SomeResponse> callCarsService(record) { return webClient.get() .uri(thirdPartyServiceUri) .accept(MediaType.APPLICATION_JSON) .retrieve() .bodyToMono(SomeResponse.class); }
3. 优化并行线程数
parallel(16)的线程数如果超过CPU核心数的合理范围(通常是核心数*2),会增加线程切换开销,降低处理效率,间接加剧数据堆积。可以根据容器CPU核心数动态调整:
.parallel(Runtime.getRuntime().availableProcessors() * 2)
为什么之前的尝试无效
- 单请求限流:每个
callCarsService的限流仅控制单个请求的发送时机,但上游已经把大量record推送到内存中等待处理,无法阻止堆积。 limitRate位置错误:如果limitRate加在parallel()之前,并行流会重新订阅上游,导致limitRate的背压控制失效;必须在并行处理后的流上配合限流和并发控制,才能有效约束整体速率。
内容的提问来源于stack exchange,提问作者some time ago
相关产品推荐
相关产品推荐

