如何在Spring WebClient中保持固定数量的并发活跃请求
问题描述
我有一个Spring Boot应用,通过Spring AMQP的@RabbitListener从RabbitMQ读取消息,消息中包含需要请求的URL,我需要请求这些URL、解析响应后保存到数据库。
我通过代理ScraperAPI发起这些URL请求,该代理最多支持50个并发请求。代理的响应耗时约30秒,因此为了让服务达到最优性能,我需要始终保持50个活跃请求。
现有代码实现
RabbitMQ消息类
public class UrlToParse{ private String propertyUrl; private String zillowUrl; }
RabbitMQ监听器
@RabbitListener(queues = "urls.queue"}) public void consumeUrls(UrlsToParse urls) { parserManager.parseUrls(urls.getUrlsToParse()); }
解析管理器
public void parseUrls(List<UrlToParse> urls) { Flux.fromIterable(stream) .parallel() .runOn(Schedulers.boundedElastic()) .flatMap(this::parse) .flatMap(url -> Mono.just(urlsRepository.save(url))) .subscribe(); } private Mono<ParsedUrl> parse(UrlToParse url) { return webClient .get() .uri(url.getUrl()) .retrieve() .bodyToMono(String.class) .flatMap(s -> Mono.just(buildParsedUrl(s))); }
我尝试过添加.delayElements()、.limitRate()以及Sinks.many().multicast().onBackpressureBuffer(),但这些看起来更像是硬编码的解决方案,当消息量增长时无法正常生效。在非响应式编程中,这个问题的解决方案通常是使用BlockingQueue,我想知道在响应式场景下有没有对应的实现机制?
最终实现方案
- 使用
flatMap(this::parse, 50)来限制flatMap的并发数 - 在Flux末尾调用
.block(),这样可以等所有URL都处理完成后再向RabbitMQ发送确认,否则确认会立即发送导致新消息被消费,之前的方式会出现同时处理10条消息、每条50并发调用,总并发达到500的问题。
如果需要更贴合响应式的方案,可以参考Reactor RabbitMQ。
替代方案(无需适配Reactor RabbitMQ)
如果不想使用Reactor RabbitMQ,可以通过Flux.generate()方法手动拉取消息,搭配flatMap(function, concurrency)来控制反压,示例代码如下:
Flux.generate((SynchronousSink<UrlToParse> synchronousSink) -> { UrlToParse urls = amqpTemplate.receiveAndConvert(queueName, UrlToParse.class); if (urls != null) { synchronousSink.next(urls); } })
内容的提问来源于stack exchange,提问作者Dmitriy Karlov
相关产品推荐
相关产品推荐

