【问题标题】:Vertx CompositeFuture: on completion of all FuturesVertx CompositeFuture:完成所有期货
【发布时间】:2022-11-04 21:32:06
【问题描述】:

在 Vert.x Web 服务器中,我有一组 Futures,每个 Futures 都可能失败或成功并保存一个结果。我对每一个 Future 的结果(可能还有结果)感兴趣,这意味着我需要处理每个 Future 的结果。

我在想 Vert.x 的 CompositeFuture 是要走的路,这是我的代码 sn-p:

List<Future> futures = dataProviders.stream()
    .filter(dp -> dp.isActive(requester))
    .map(DataProvider::getData)
    .collect(Collectors.toList());

CompositeFuture.all(futures)
        .onComplete(ar -> {
            if(ar.failed()) {
                routingContext.response()
                    .end(ar.cause());
                return;
            }
            
            CompositeFuture cf = ar.result();
            JsonArray data = new JsonArray();
            for(int i = 0; i < cf.size(); i++) {
                if(cf.failed(i)) {
                    final JsonObject errorJson = new JsonObject();
                    errorJson.put("error", cf.cause(i).getMessage());
                    data.add(errorJson);
                } else {
                    data.add(((Data) cf.resultAt(i)).toJson());
                }
            }

            JsonObject res = new JsonObject()
                .put("data", data);

            routingContext.response()
                    .putHeader("Content-Type", "application/json")
                    .end(res.toString());
        });

但随之而来的是以下问题:

  • 使用CompositeFuture.all(futures).onComplete(),只要futures 中的任何Future 失败(因为ar.result() 为空),我就不会得到成功Future 的结果。
  • 使用 CompositeFuture.any(futures).onComplete(),我会得到所有结果,但 CompositeFuture 会在 futures 的所有 Future 完成之前完成。意思是,它不会等待每个 Future 完成,而是在任何 Future 完成后立即完成。 (-> cf.resultAt(i) 返回空值)
  • 使用CompositeFuture.join(futures).onComplete(),与all() 相同:只要任何Future 失败,ar.result() 就会为空。

等待 Futures 列表完成,同时能够单独处理每个结果和结果的正确/最佳方法是什么?

【问题讨论】:

    标签: java asynchronous vert.x


    【解决方案1】:

    最简单的方法是您自己处理结果。您可以将onSuccess 处理程序注册到您的期货。这样,结果将被放入某种列表中,例如JsonArray

    List<Future> futures = //...list of futures
    
    JsonArray results = new JsonArray();
    futures.forEach(e -> e.onSuccess(h -> results.add(h)));
    
    CompositeFuture.all(futures)
        .onComplete(ar -> {
            if(ar.failed()) {
                // successful elements are present in "results"
                routingContext.response().end(results.encode());
                return;
            }
            //... rest of your code
         });
    

    您还可以查看rx-java 库。使用它通常可以更好地实现此类用例。

    【讨论】:

    • 谢谢!也应该知道这一点......只添加了一个计数器以保持结果与之前的期货相同(否则在完成时间之后排序)。
    • 不客气,这很好,订单不能保证。
    【解决方案2】:

    使用all 时,您可以简单地戳原始期货的结果:

    List<Future> futures = //...list of futures
    
    CompositeFuture.all(futures).onComplete(ar -> {
      if(ar.succeeded()){
        futures.forEach(fut -> log.info( fut.succeded() +" / " +_fut.result() ));
      }
    } );
    

    【讨论】:

      【解决方案3】:

      我会使用join 操作而不是全部操作,因此即使其中一个失败,CompositeFuture 也会等待所有期货完成。否则,与 CompositeFuture 的评估绑定可能会在所有异步调用完成之前发生。

      然后,您可以在 CompositeFuture 中加入期货之前定义recover 步骤。

        public static void main(String[] args) {
          List<String> inputs = List.of("1", "Terrible error", "2");
      
          List<Future> futures = inputs.stream()
              .map(i -> asyncAction(i)
                  .recover(thr -> {
                    // real recovery or just logging
                    System.out.println("Bad thing happen: " + thr.getMessage());
                    return Future.succeededFuture();
                  }))
              .collect(Collectors.toList());
      
          CompositeFuture.join(futures)
              .map(CompositeFuture::list)
              .map(results -> results.stream()
                  // filter out empty recovered future            
                  .filter(Objects::nonNull)
                  .collect(Collectors.toList()))
              .onSuccess(System.out::println);
        }
      
        static Future<String> asyncAction(final String input) {
          if ("Terrible error".equals(input)) {
            return Future.failedFuture(input);
          }
          return Future.succeededFuture(input);
        }
      

      它将打印:

      Bad thing happen: Terrible error
      [1, 2]
      

      【讨论】:

        猜你喜欢
        • 2015-05-17
        • 1970-01-01
        • 2021-12-30
        • 1970-01-01
        • 2021-12-29
        • 2023-03-16
        • 1970-01-01
        • 2023-03-03
        • 2013-06-29
        相关资源
        最近更新 更多