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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 02:57:03