【问题标题】:How to create a Source from Future[Iterator]?如何从 Future[Iterator] 创建源?
【发布时间】:2018-05-09 19:57:39
【问题描述】:

我有一个Future[Iterator]。我想将此迭代器提供给我的 Stream。在这里,我仍然想像使用 Source.fromIterator 一样从 Iterator 构造一个 Source。

但是,由于 Future,我不能在这里使用 Source.fromIterator

也许我可以使用Source.fromFuture,但是当我尝试使用它时,在我的情况下,它似乎实际上并没有从迭代器创建一个源。来自文档:

/**
   * Starts a new `Source` from the given `Future`. The stream will consist of
   * one element when the `Future` is completed with a successful value, which
   * may happen before or after materializing the `Flow`.
   * The stream terminates with a failure if the `Future` is completed with a failure.
   */
  def fromFuture[T](future: Future[T]): Source[T, NotUsed] =
    fromGraph(new FutureSource(future))

【问题讨论】:

    标签: scala stream akka akka-stream


    【解决方案1】:

    您可以使用Source.fromFutureflatMapConcatSource.fromIterator 的组合。例如:

    val futIter = Future(Iterator(1, 2, 3))
    
    val source: Source[Int, _] =
      Source.fromFuture(futIter).flatMapConcat(iter => Source.fromIterator(() => iter))
    

    【讨论】:

    • 加 1,以获得出色的提示。不知道这个方法。
    【解决方案2】:

    Chunjef 的回答最符合问题的要求。

    我只是想指出,这种功能通常是通过mapflatMap 在另一个Future 中完成的:

    type Data = ???
    
    val iterFuture : Future[Iterator[Data]] = ???
    
    val dataSeq : Future[Seq[Data]] = iterFuture flatMap { iter =>
      Source
        .fromIterator(() => iter)
        .to(Sink.seq[Data])
        .run()
    }
    

    如果您根本不想使用流,还有 Future.traverse 函数:

    type OutputData = ???
    
    val someCalculation : Data => Future[OutputData] = ???
    
    val outputIterFuture : Future[Iterator[OutputData]] = 
      iterFuture flatMap { iter => Future.traverse(iter)(someCalculation) }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-03-24
      • 2021-06-07
      • 1970-01-01
      • 2012-06-16
      • 1970-01-01
      • 2021-09-05
      • 2011-06-16
      相关资源
      最近更新 更多