【问题标题】:How to combine Source.repeat and Source.completionStage using Akka如何使用 Akka 结合 Source.repeat 和 Source.completionStage
【发布时间】:2021-04-27 10:23:37
【问题描述】:

我正在使用带有微服务框架的 akka,所以我收到了很多完成阶段请求。我想从一个微服务中获取元素列表,并将它们与另一个微服务中的单个元素一起压缩,这样我最终得到一个 Source of Pair

我无法使用普通 zip 执行此操作,因为 Source.zip 会在两个源之一完成后立即完成,因此我最终只会传播一个元素。

我不能使用 Source.zipAll,因为这需要我提前定义默认元素。

如果我提前拥有单个元素,我可以使用 Source.repeat 使其重复传播该元素,这意味着 Source.zip 将在元素列表完成时完成,但 Source.repeat 可以t 采取完成阶段或 Source.completionStage。

我目前的策略是在 mapConcat 列表元素之前将所有内容压缩在一起。

Source<singleElement> singleElement = Source.completionStage(oneService.getSingleElement().invoke());

return Source.completionStage(anotherService.getListOfElements().invoke)
    .zip(singleElement)
    .flatMapConcat(pair -> Source.fromIterator(() -> pair.first().stream().map(listElement -> Pair.create(listElement, pair.second())));

这最终得到了我想要的,但我觉得有很多不必要的重复和同步移动数据。有没有更好的方法来解决我错过的这个问题?

【问题讨论】:

    标签: java akka akka-stream lagom completion-stage


    【解决方案1】:

    flatMapConcat 运算符应该允许您构造一个 Source.repeat,它会在已知单个元素后重复该元素。在 Scala 中(Source.future 相当于 Source.completionStage 的 Scala:我对 Java lambda 语法不够熟悉,无法用 Java 回答):

    val singleElement = Source.future(oneService.getSingleElement)
    
    Source.future(anotherService.getListOfElements)
      .mapConcat(lst => lst)  // unspool the list
      .zip(singleElement.flatMapConcat(element => Source.repeat(element)))
    

    【讨论】:

      【解决方案2】:

      为什么不将CompletionStages 组合起来,然后将它们提供给 Akka 流?

      Source<Pair<String,String>, ?> execute() {
          CompletionStage<Pair<String, List<String>>> pairCompletionStage = getSingleElement().thenCombine(getListOfElements(), Pair::create);
      
          return Source.completionStage(pairCompletionStage)
              .flatMapConcat(pair -> Source.from(pair.second()).map(listElement -> Pair.create(listElement, pair.first())));
      }
      

      完整的 PoC - 玩睡眠超时以先完成一个或另一个 CompletionStage

      import java.util.Arrays;
      import java.util.List;
      import java.util.concurrent.CompletableFuture;
      import java.util.concurrent.CompletionStage;
      
      import akka.Done;
      import akka.actor.ActorSystem;
      import akka.japi.Pair;
      import akka.stream.javadsl.Sink;
      import akka.stream.javadsl.Source;
      
      public class CompletionStages {
          CompletionStage<String> getSingleElement() {
              return CompletableFuture.supplyAsync(() -> {
                  try {
                      Thread.sleep(5000);
                      return "Single Element";
                  } catch (InterruptedException e) {
                      Thread.currentThread().interrupt();
                      return null;
                  }
              });
          }
      
          CompletionStage<List<String>> getListOfElements() {
              return CompletableFuture.supplyAsync(() -> {
                  try {
                      Thread.sleep(3000);
                      return Arrays.asList("One", "Two", "Three");
                  } catch (InterruptedException e) {
                      Thread.currentThread().interrupt();
                      return null;
                  }
              });
          }
      
          Source<Pair<String,String>, ?> execute() {
              CompletionStage<Pair<String, List<String>>> pairCompletionStage = getSingleElement().thenCombine(getListOfElements(), Pair::create);
      
              return Source.completionStage(pairCompletionStage)
                      .flatMapConcat(pair -> Source.from(pair.second()).map(listElement -> Pair.create(listElement, pair.first())));
          }
      
          CompletionStage<Done> run(ActorSystem system) {
              return execute().runWith(Sink.foreach(System.out::println), system);
          }
      
          public static void main(String... args) {
              final ActorSystem system = ActorSystem.create();
              new CompletionStages().run(system)
                      .thenRun(system::terminate);
          }
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-09-02
        • 2011-10-28
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2014-02-18
        • 2019-08-02
        相关资源
        最近更新 更多