【问题标题】:Flux from WebClient behaves differently than Flux from File.readLines来自 WebClient 的 Flux 的行为与来自 File.readLines 的 Flux 不同
【发布时间】:2019-05-01 13:58:27
【问题描述】:

我必须使用两个不同的 ID 来源。一个来自文件,另一个来自 URL。当我从文件的行创建Flux 时,我可以很好地处理它。当我将 Flux-creating 函数与使用 WebClient....get() 的函数切换时,我得到了不同的结果;由于某种原因,WebClient 永远不会被调用。

private Flux<String> retrieveIdListFromFile(String filename) {
  try {
    return Flux.fromIterable(Files.readAllLines(ResourceUtils.getFile(filename).toPath()));
  } catch (IOException e) {
    return Flux.error(e);
  }
}

这里是 WebClient 部分...

private Flux<String> retrieveIdList() {
  return client.get()
      .uri(uriBuilder -> uriBuilder.path("capdocuments_201811v2/selectRaw")
          .queryParam("q", "-P_Id:[* TO *]")
          .queryParam("fq", "DateLastModified:[2010-01-01T00:00:00Z TO 2016-12-31T00:00:00Z]")
          .queryParam("fl", "id")
          .queryParam("rows", "10")
          .queryParam("wt", "csv")
          .build())
      .retrieve()
      .bodyToFlux(String.class);
}

当我在 WebClient 的通量上执行 subscribe(System.out::println) 时,没有任何反应。当我执行 blockLast() 时,它可以工作(调用 URL,返回数据)。我不明白为什么,如何纠正,以及我做错了什么。 使用源自文件的通量,即使订阅也可以正常工作。我有点想,通量是可以互换的......

当我做retrieveIdList().log().subscribe():

INFO [main] reactor.Flux.OnAssembly.1 | onSubscribe([Fuseable] FluxOnAssembly.OnAssemblySubscriber)
INFO [main] reactor.Flux.OnAssembly.1 | request(unbounded)

当我用 blockLast() 而不是 subscribe() 做同样的事情时:

INFO [main] reactor.Flux.OnAssembly.1 | onSubscribe([Fuseable] FluxOnAssembly.OnAssemblySubscriber)
INFO [main] reactor.Flux.OnAssembly.1 | request(unbounded)
INFO [reactor-http-nio-4] reactor.Flux.OnAssembly.1 | onNext(id)
.
.
.

【问题讨论】:

  • 你能在 webclient 上添加一个log() 运算符来看看发生了什么吗?另外,是否可以在底层 reactor netty HttpClient 上启用wiretap,以便我们可以看到 HTTP 请求/响应?
  • 您能用所要求的信息更新您的问题吗?您还缺少 HTTP 日志。
  • 什么都没有,没有发送请求。

标签: spring-webflux project-reactor


【解决方案1】:

从您的问题更新来看,似乎没有什么正在等待处理完成。我假设这是批处理或 CLI 应用程序,而不是 Web 应用程序?

假设如下:

Flux<User> users = userService.fetchAll();

Flux 上调用blockLast 将触发处理和block,直到结果出现。

在它上面调用subscribe会触发异步处理;我们在您的日志中看到了订阅者 request 元素,但仅此而已。这可能意味着 JVM 在发布任何元素之前退出 - 没有任何东西在等待结果。

如果您正在有效地编写一些 CLI/批处理应用程序而不是在 Web 应用程序中处理请求,您可以在最终的反应管道上block 以获得结果。如果您希望将该结果写入文件或将其发送到不同的服务,那么您应该使用反应器操作员对其进行撰写。

【讨论】:

  • 啊哈,那很可能......我确实将它用作批处理,抱歉我没有提到!我现在的问题是:我如何并行化它? ParallelFlux 上没有块功能,对吧?您是否有任何指示?我确定这已经完成了;-)
  • 并行化是什么意思?在您收集结果之前,一切都以非阻塞方式异步完成。
  • 使用parallel() 和runOn() 不会让你blockLast(),所以我现在明白我的问题表述有误......这里是阻塞的解决方案- .retrieveIdList() .parallel( 10) .log() .runOn(Schedulers.parallel()) .sequential() //
  • 这是一个完全不同的问题,我不明白。您能否提出一个不同的问题并提供一些有关您要实现的目标的背景,展示您如何尝试实现它的代码 sn-p,预期结果是什么以及您得到的结果是什么?
猜你喜欢
  • 2021-09-11
  • 2016-02-27
  • 2020-10-08
  • 2019-09-21
  • 2023-02-05
  • 2016-03-03
  • 2022-11-30
  • 2022-01-08
  • 1970-01-01
相关资源
最近更新 更多