【问题标题】:asynchronously prepend to stream异步添加到流
【发布时间】:2015-07-15 23:08:22
【问题描述】:

我有一个函数,它返回一个对源自套接字的无限流执行操作。

def f1(s:Stream[T]):Stream[T] = s map {...} filter {...}

我有另一个方法可以返回一个有限序列,我想加入这个流。

def f2():Stream[T] = ...

这就是我要找的东西:

val a = f2#:::f1(s)

问题是,f2 需要一些时间来计算,而f1 需要尽快输出值。我想将f2 包装在Runnable 中,这样我就可以在不阻塞程序其余部分的情况下计算它的值。

我希望 a 的行为如下:如果 f2 已完成计算,请将其添加到 f1 并继续输出该流的值。否则,继续输出f1 的值。我应该如何进行?

【问题讨论】:

  • 如果f2没有计算完,我们继续输出f1的值,然后f2完成计算会怎样?
  • 完成后添加到 f1 -- 所以 steam 的定义是 f2 while any f2 else f1

标签: scala asynchronous stream


【解决方案1】:

您可以构建一个扩展Stream 的自定义类,并为您想要的语义定义所有必要的方法。在定义自定义语义时,采用这种方法将为您提供最大的灵活性。您甚至可以丰富Stream,使其具有withAsynPrefix 方法。下面是Iterator的粗略版本,稍微简单一点:

import scala.concurrent._
import scala.concurrent.duration._
import scala.util._
import scala.concurrent.ExecutionContext.Implicits.global

case class CombinedIterator[T](val iterator: Iterator[T], val prefix: Future[Iterator[T]]) extends Iterator[T] {
    def hasNext: Boolean = iterator.hasNext || Await.result(prefix, Duration.Inf).hasNext
    def next: T = prefix.value.collect {
        case Success(pfx) if pfx.hasNext => pfx.next
    }.getOrElse(iterator.next)
}

implicit class AsyncPrefixable[T](val iterator: Iterator[T]) extends AnyVal {
    def withAsyncPrefix(prefix: Future[Iterator[T]]): Iterator[T] = CombinedIterator(iterator, prefix)
}

那么你可以这样做:

val iterator = Iterator.from(0).take(10000)
val prefix = Future { Thread.sleep(10); Iterator(-1) }
val iteratorWithPrefix = iterator.withAsyncPrefix(prefix)
iteratorWithPrefix.toList.indexOf(-1) //Run a few times, location can vary greatly within list

【讨论】:

  • 本,这看起来很棒。我可以在这个Iterator 上调用.toStream 来生成Stream 吗?
  • @WalrustheCat 可能!您也可以扩展Stream 而不是Iterator。我根本没有测试过这个,所以你应该确保它捕捉到任何边缘情况!
猜你喜欢
  • 2020-01-24
  • 2015-08-12
  • 1970-01-01
  • 2019-01-02
  • 2011-11-30
  • 2015-08-18
  • 2020-01-02
  • 1970-01-01
  • 2016-11-22
相关资源
最近更新 更多