【问题标题】:Mono returned by ServerRequest.bodyToMono() method not extracting the body if I return ServerResponse immediately如果我立即返回 ServerResponse,ServerRequest.bodyToMono() 方法返回的 Mono 不会提取正文
【发布时间】:2020-01-05 16:29:19
【问题描述】:

我在 Spring Web Flux 中使用 Web Reactive。我已经为 POST 请求实现了一个 Handler 函数。我希望服务器立即返回。所以,我实现了如下处理程序 - :

public class Sample implements HandlerFunction<ServerResponse>{

public Mono<ServerResponse> handle(ServerRequest request) {

Mono bodyMono = request.bodyToMono(String.class);

bodyMono.map(str -> {
  System.out.println("body got is " + str);
  return str;
}).subscribe();

return ServerResponse.status(HttpStatus.CREATED).build();
    }
}

但是 map 函数中的 print 语句没有被调用。这意味着身体没有被提取出来。 如果我不立即返回响应并使用

return bodyMono.then(ServerResponse.status(HttpStatus.CREATED).build())

然后地图函数被调用。

那么,如何在后台处理我的请求正文? 请帮忙。

编辑

我尝试使用下面的flux.share() -:

Flux<String> bodyFlux = request.bodyToMono(String.class).flux().share();
Flux<String> processFlux = bodyFlux.map(str -> {
      System.out.println("body got is");
      try{
        Thread.sleep(1000);
      }catch (Exception ex){

      }
      return str;
    });

    processFlux.subscribeOn(Schedulers.elastic()).subscribe();

    return bodyFlux.then(ServerResponse.status(HttpStatus.CREATED).build());

在上面的代码中,map 函数有时被调用,有时不被调用。

【问题讨论】:

  • 您发布的代码无法编译。发布实际的编译代码,重现问题。
  • 我已编辑问题以发布全班。
  • 而且代码仍然无法编译。
  • 现在尝试编译

标签: spring spring-webflux project-reactor


【解决方案1】:

这是不可能的。

Web 服务器(包括 Reactor Netty、Tomcat 等)在请求处理完成时清理和回收资源。这意味着当您的控制器处理程序完成时,HTTP 资源、请求本身、可重用缓冲区等将被回收或关闭。此时,您无法再读取请求正文。

在您的情况下,您需要先读取并缓冲整个请求正文,然后返回响应并启动一个任务以在单独的执行中处理该请求。

【讨论】:

  • 如何先读取和缓冲整个请求体?我不能使用 block(),因为它是 nio 线程。
【解决方案2】:

正如您所发现的,您不能随意将subscribe() 分配给bodyToMono() 返回的Mono,因为在这种情况下,主体根本不会传递到Mono 进行处理。 (您可以通过在该Mono 中调用single() 来验证这一点,它会抛出异常,因为不会发出任何元素。)

那么,如何在后台处理我的请求正文?

如果你真的仍然想只是使用reactor在后台做一个很长的任务同时立即返回,你可以这样做:

return request.bodyToMono(String.class).doOnNext(str -> {
    Mono.just(str).publishOn(Schedulers.elastic()).subscribe(s -> {
        System.out.println("proc start!");
        try {
            Thread.sleep(1000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println("proc end!");
    });
}).then(ServerResponse.status(HttpStatus.CREATED).build());

这种方法立即将发出的元素发布到新的Mono,设置为在弹性调度程序上发布,然后在后台订阅。然而,它有点难看,而且这并不是反应堆的设计目的。您可能在这里误解了反应器/反应式编程背后的想法:

  • 它不是以“返回快速结果然后在后台执行操作”的想法编写的——这通常是工作队列的目的,通常使用 RabbitMQ 或 Kafka 之类的东西来实现。它的“存在理由”是非阻塞,因此单个线程永远不会被空闲阻塞,等待其他事情完成。
  • map() 方法不是为副作用而设计的,它旨在将每个对象转换为另一个对象。对于副作用,您需要doOnNext();
  • Reactor 默认使用单个线程,因此您在 map() 方法中的“附加处理”仍会阻塞该线程。

如果您的应用程序不仅仅用于快速演示目的,并且/或者您需要大量使用此模式,那么我会认真考虑设置一个适当的工作队列。

【讨论】:

  • 我还使用 share() 添加了另一种方法。在这种情况下,有时会调用 map 函数,有时则不会。
猜你喜欢
  • 2021-09-09
  • 2017-12-11
  • 1970-01-01
  • 2013-05-12
  • 1970-01-01
  • 1970-01-01
  • 2021-03-25
  • 2012-07-11
  • 2013-03-08
相关资源
最近更新 更多