【问题标题】:Parallel GET Request to specific mapping with WebFlux使用 WebFlux 对特定映射的并行 GET 请求
【发布时间】:2020-06-13 05:13:48
【问题描述】:

我想与WebClient 同时调用独立请求。我之前使用RestTemplate 的方法在等待响应时阻塞了我的线程。所以我发现,WebClient 和 ParallelFlux 可以使用一个线程更高效,因为它应该用一个线程调度多个请求。

我的端点请求一个 id 和一个 location 的元组。

fooFlux 方法将在具有不同参数的循环中被调用数千次。返回的地图将根据存储的参考值进行断言。

之前对该代码的尝试导致重复的 API 调用。 但是仍然有一个缺陷。 mapping 的键集大小通常小于Set<String> location 的大小。事实上,生成的地图的大小正在发生变化。此外,它时不时地是正确的。因此,在方法返回地图后,下标完成可能会出现问题。

public Map<String, ServiceDescription> fooFlux(String id, Set<String> locations) {
    Map<String, ServiceDescription> mapping = new HashMap<>();
    Flux.fromIterable(locations).parallel().runOn(Schedulers.boundedElastic()).flatMap(location -> {
        Mono<ServiceDescription> sdMono = getServiceDescription(id, location);
        Mono<Mono<ServiceDescription>> sdMonoMono = sdMono.flatMap(item -> {
            mapping.put(location, item);
            return Mono.just(sdMono);
        });
        return sdMonoMono;
    }).then().block();
    LOGGER.debug("Input Location size: {}", locations.size());
    LOGGER.debug("Output Location in map: {}", mapping.keySet().size());
    return mapping;
}

处理获取请求

private Mono<ServiceDescription> getServiceDescription(String id, String location) {
    String uri = URL_BASE.concat(location).concat("/detail?q=").concat(id);
    Mono<ServiceDescription> serviceDescription =
                    webClient.get().uri(uri).retrieve().onStatus(HttpStatus::isError, clientResponse -> {
                        LOGGER.error("Error while calling endpoint {} with status code {}", uri,
                                        clientResponse.statusCode());
                        throw new RuntimeException("Error while calling Endpoint");
                    }).bodyToMono(ServiceDescription.class).retryBackoff(5, Duration.ofSeconds(15));
    return serviceDescription;
}

【问题讨论】:

  • 为什么你使用JsonNode.class而不是序列化/反序列化成一个具体的对象?以及为什么使用反应式编程来解决可以使用@Async 解决的问题。反应式编程不是异步编程。它们是相辅相成的两种不同的东西。
  • 我使用了JsonNode.class,因为收到的 JSON 模型很大,我只需要它的一小部分。由于一篇 baeldung 文章 (baeldung.com/spring-webclient-resttemplate),我提出了响应式编程。我想存档下载速度的提升。 ParallelStreams 中的RestTemplate 方法让我的网卡下载速度在 10Mbit/s 到 200Mbit/s 之间。取决于每个 id 的位置数量。但这从 1 到 ~4000 不等
  • RestClients 不会影响下载速度。网络带宽影响速度。所以你使用什么客户端不会影响任何下载速度。如果你只需要一小部分,谁说你需要声明整个对象?只需使用您需要的小块创建一个类。是的,使用 WebClient,但你需要知道反应式编程和并发编程之间的区别。
  • 响应式编程就是不阻塞,并尽可能多地利用线程来执行串行和并行任务。虽然并发编程是做你想做的事,但同时产生线程并获取东西。反应式编程可以串行和并行地做事情来解决你给他们的任务。但它并不主要用于执行异步任务。
  • 好吧,如果您需要收集所有结果并获取具体值以在一个大块中返回到调用客户端,而不是将结果流式传输到调用客户端,那么在您的应用程序中不是响应式的, 需要使用块。但是正如您所做的那样,放置在自己的调度程序上。但我建议使用 boundedElastic 调度程序,您选择的调度程序会将线程耗尽到无穷大,在最坏的情况下会遇到线程不足并导致应用程序崩溃。

标签: java spring-webflux project-reactor spring-webclient


【解决方案1】:
public Map<String, ServiceDescription> fooFlux(String id, Set<String> locations) {
    return Flux.fromIterable(locations)
               .flatMap(location -> getServiceDescription(id, location).map(sd -> Tuples.of(location, sd)))
               .collectMap(Tuple2::getT1, Tuple2::getT2)
               .block();
}

注意:flatMap 运算符与 WebClient 调用结合使用可以并发执行,因此无需使用 ParallelFlux 或任何 Scheduler。

【讨论】:

  • 谢谢马丁,我认为这就是重点。在那种情况下,我无法找到或正确使用 collectMap 方法。您能否详细说明为什么不使用ParallelFlux 和Scheduler?我认为工作单元的共享线程池会在 I/O 情况下加速。
  • 据我了解,ParallelFlux 用于 CPU 密集型工作。没有它也可以实现并发 IO,因为在等待外部资源期间不使用 CPU,因此可以用非常少的线程数实现高并发。 ParallelFlux 只是让事情变得比需要的更复杂。
【解决方案2】:

当您订阅生产者时,响应式代码会被执行。 Block 确实订阅了,因为您调用了两次 block(一次在 Mono 上,但再次返回 Mono,然后在 ParallelFlux 上调用 block),Mono 被执行两次。

    List<String> resultList = listMono.block();
    mapping.put(location, resultList);
    return listMono;

尝试以下类似的方法(未经测试):

    listMono.map(resultList -> {
       mapping.put(location, resultList);
       return Mono.just(listMono);
    });

也就是说,响应式编程模型非常复杂,因此请考虑使用 @Async 和 Future/AsyncResult,如果这只是关于并行调用远程调用,正如其他人所建议的那样。 您仍然可以使用 WebClient(RestTemplate 似乎即将被弃用),但只需在 bodyToMono 之后立即调用 block。

【讨论】:

  • 感谢您的回复。这有助于我消除重复的请求。为了防止出现空地图,我必须将Mono.just(listMono) vlaue 绑定到我需要在外部 flatMap 末尾返回的变量。起初我试图返回listMono 而不用Mono.just(listMono) 创建一个新的Mono。这也导致了重复的请求。你能解释一下这种行为吗?
  • 不幸的是,填充的地图没有填充完整。键集mapping 的大小通常小于Set&lt;String&gt; location 的大小。事实上,生成的地图的大小正在发生变化。我是否需要内部订阅或其他东西来确保在方法返回之前完全填充地图?我更新了我的问题,以反映我在此处回答期间所做的更改。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-02
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多