【问题标题】:Keep doing a job after client timeouts / return current progress before timeout客户端超时后继续工作/超时前返回当前进度
【发布时间】:2019-08-15 12:38:10
【问题描述】:

我正在编写一个适用于下载和持久化数据的 Spring-Webflux 应用程序。

一个用例是来自客户端的 REST 请求,用于触发将数据从外部依赖项下载到应用程序的数据库并返回下载的总条目。

我将此下载作为对Flux 的订阅处理,并以包含告诉用户下载了多少数据的json 的Mono 进行响应:

http :8080/rest/download from==2019-08-14T00:00:00Z

HTTP/1.1 200 OK
Content-Length: 15
Content-Type: application/json;charset=UTF-8

{
    "count": 19000
}

使用的 Java 代码如下所示:

public Mono<Map<String, Object>> importDataBetween(Instant from, Instant to) {
    return backendAdapter.importData(from, to)
            .doOnComplete(this::cleanupDb)
            .doOnNext(this::storeData)
            .count()
            .map(count -> Map.of("count", count)); // e.g. { 'count' : 123 }
}

现在的问题是,如果客户端超时,作业本身就会被取消:

http :8080/rest/download from==2019-08-01T00:00:00Z

http: error: Request timed out (30s).

webflux 是否有任何标准功能来处理此类情况?

是否可以计算结果并在超时后计算,例如29s 返回那个号码并附上“下载仍在进行中”的备注?

理想情况下,如果操作时间少于 30 秒,则应返回结果,否则可以返回当前进度或默认响应,但不应取消对下载过程的订阅。

非常感谢!

【问题讨论】:

    标签: spring-webflux project-reactor


    【解决方案1】:

    有几种方法可以解决您的一些问题。

    你可以使用Flux#bufferTimeout(int maxSize, Duration maxTime)

    要缓冲多个项目,请在项目数量达到 maxsize 或达到以秒为单位的时间时结束发射。

    有几种方法,我建议查看 Flux api 超时和缓冲区部分以找到适合您的方法。

    Flux API

    【讨论】:

      【解决方案2】:

      有两种可能性:如果您想在客户端估计超时之前提前关闭连接,您可以使用具有合理持续时间的Flux::take。如果您想在响应中添加任意信息,请使用Flux::buffer 并将自定义消息添加到已缓冲的对象中。

      为了展示这一点,请考虑以下示例:

      @Slf4j
      @RestController
      @RequestMapping("api")
      public class DemoController {
      
          @GetMapping("flux-take")
          public Flux<Result> fluxTake() {
              Flux<Result> resultFlux = createFlux();
              resultFlux.share().subscribe(result -> log.info("result={}", result));
              return resultFlux.take(Duration.ofMillis(1500));
          }
      
          @GetMapping("flux-buffer")
          public Mono<BufferedResult> fluxBuffer() {
              Flux<Result> resultFlux = createFlux();
              resultFlux.share().subscribe(result -> log.info("result={}", result));
              return resultFlux.buffer(Duration.ofMillis(1500))
                      .next()
                      .map(results -> new BufferedResult("Request timed out, returning buffered results", results));
          }
      
          private Flux<Result> createFlux() {
              return Flux.just(1, 2, 3, 4, 5, 6)
                      .delayElements(Duration.ofMillis(300))
                      .timed()
                      .map(timedId -> new Result(timedId.get(), timedId.elapsed().toMillis()));
          }
      
          @Data
          @AllArgsConstructor
          public static class Result {
              private final Integer id;
              private final Long elapsedMillis;
          }
      
          @Data
          @AllArgsConstructor
          public static class BufferedResult {
              private final String message;
              private final List<Result> results;
          }
      }
      

      从客户的角度来看,结果是这样的:

      $ http :8080/api/flux-take --timeout=2
      HTTP/1.1 200 OK
      Content-Type: application/json
      transfer-encoding: chunked
      
      [
          {
              "elapsedMillis": 304,
              "id": 1
          },
          {
              "elapsedMillis": 300,
              "id": 2
          },
          {
              "elapsedMillis": 304,
              "id": 3
          },
          {
              "elapsedMillis": 302,
              "id": 4
          }
      ]
      
      $ http :8080/api/flux-buffer --timeout=2
      HTTP/1.1 200 OK
      Content-Length: 187
      Content-Type: application/json
      
      {
          "message": "Request timed out, returning buffered results",
          "results": [
              {
                  "elapsedMillis": 301,
                  "id": 1
              },
              {
                  "elapsedMillis": 301,
                  "id": 2
              },
              {
                  "elapsedMillis": 302,
                  "id": 3
              },
              {
                  "elapsedMillis": 304,
                  "id": 4
              }
          ]
      }
      

      请注意客户端的 2 秒超时和应用程序内的 1500 毫秒超时。客户端永远不会达到 2 秒超时,并且总是会收到正确的响应。

      IntelliJ 将在 resultFlux.share().subscribe(...) 上显示提示 Calling 'subscribe' in non-blocking scope,但我想不出更好的方法来完成此任务。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-10-10
        • 1970-01-01
        • 2019-09-22
        • 2023-03-19
        • 1970-01-01
        • 2013-12-05
        • 1970-01-01
        • 2019-05-05
        相关资源
        最近更新 更多