【问题标题】:restart iterator on exceptions in Scala在Scala中的异常上重新启动迭代器
【发布时间】:2010-10-24 20:27:14
【问题描述】:

我有一个迭代器(实际上是 Source.getLines),它从 URL 读取无限的数据流。当出现连接问题时,迭代器偶尔会抛出java.io.IOException。在这种情况下,我需要重新连接并重新启动迭代器。我希望这是无缝的,这样迭代器对消费者来说就像一个普通的迭代器,但在必要时会在下面重新启动。

例如,我希望看到以下行为:

scala> val iter = restartingIterator(() => new Iterator[Int]{
  var i = -1
  def hasNext = {
    if (this.i < 3) {
      true
    } else {
      throw new IOException
    }
  }
  def next = {
    this.i += 1
    i
  }
})
res0: ...

scala> iter.take(6).toList
res1: List[Int] = List(0, 1, 2, 3, 0, 1)

我对这个问题有部分解决方案,但它会在某些极端情况下失败(例如,重新启动后第一个项目上的 IOException)并且非常难看:

def restartingIterator[T](getIter: () => Iterator[T]) = new Iterator[T] {
  var iter = getIter()
  def hasNext = {
    try {
      iter.hasNext
    } catch {
      case e: IOException => {
        this.iter = getIter()
        iter.hasNext
      }
    }
  }
  def next = {
    try {
      iter.next
    } catch {
      case e: IOException => {
        this.iter = getIter()
        iter.next
      }
    }
  }
}

我一直觉得有更好的解决方案,可能是 Iterator.continuallyutil.control.Exception 的某种组合或类似的东西,但我想不出一个。有什么想法吗?

【问题讨论】:

  • 我在原始答案中添加了带有continuallyutil.control.Exception 的解决方案。

标签: exception scala iterator restart


【解决方案1】:

这与您的版本相当接近,并使用scala.util.control.Exception

def restartingIterator[T](getIter: () => Iterator[T]) = new Iterator[T] {
  import util.control.Exception.allCatch
  private[this] var i = getIter()
  private[this] def replace() = i = getIter()
  def hasNext: Boolean = allCatch.opt(i.hasNext).getOrElse{replace(); hasNext}
  def next(): T = allCatch.opt(i.next).getOrElse{replace(); next}
}

由于某种原因,这不是尾递归,但可以通过使用稍微详细一点的版本来解决:

def restartingIterator2[T](getIter: () => Iterator[T]) = new Iterator[T] {
  import util.control.Exception.allCatch
  private[this] var i = getIter()
  private[this] def replace() = i = getIter()
  @annotation.tailrec def hasNext: Boolean = {
    val v = allCatch.opt(i.hasNext)
    if (v.isDefined) v.get else {replace(); hasNext}
  }
  @annotation.tailrec def next(): T = {
    val v = allCatch.opt(i.next)
    if (v.isDefined) v.get else {replace(); next}
  }
}

编辑:有util.control.ExceptionIterator.continually的解决方案:

def restartingIterator[T](getIter: () => Iterator[T]) = {
  import util.control.Exception.allCatch
  var iter = getIter()
  def f: T = allCatch.opt(iter.next).getOrElse{iter = getIter(); f}
  Iterator.continually { f }
}

【讨论】:

  • 是的,使其递归解决了我有点担心的极端情况。我想我可以通过将解决方案中的第二个“iter.hasNext”和“iter.next”更改为“this.hasNext”和“this.next”并添加 talrec 注释来获得几乎相同的行为。不过,我有点希望有一个基于组合的更简单的解决方案。
  • 非常酷。这正是我一直在寻找的东西,谢谢!
  • @huynhji- 我对 sn-ps if (v.isDefined) v.get else {replace(); 有点困惑next} 和 if (v.isDefined) v.get else {replace();有下一个}。如果发生异常,这两行不要将迭代器一直重置到开始。我试图了解它如何跳过引发异常的部分并移动到它正在迭代的源的下一个元素?
  • @sc_ray,试图记住(这是很久以前的事了),它确实会重新启动迭代器,但这就是问题所在。
【解决方案2】:

有一个更好的解决方案,Iteratee:

http://apocalisp.wordpress.com/2010/10/17/scalaz-tutorial-enumeration-based-io-with-iteratees/

这里是一个枚举器,遇到异常会重新启动。

def enumReader[A](r: => BufferedReader, it: IterV[String, A]): IO[IterV[String, A]] = {
  val tmpReader = r
  def loop: IterV[String, A] => IO[IterV[String, A]] = {
    case i@Done(_, _) => IO { i }
    case Cont(k) => for {
      s <- IO { try { val x = tmpReader.readLine; IO(x) }
                catch { case e => enumReader(r, it) }}.join
      a <- if (s == null) k(EOF) else loop(k(El(s)))
    } yield a
  }
  loop(it)
}

内部循环推进 Iteratee,但外部函数仍保留原始函数。由于Iteratee是一个持久化的数据结构,要重新启动你只需要再次调用该函数。

我在这里按名称传递 Reader,以便 r 本质上是一个为您提供全新(重新启动)阅读器的函数。在实践中,您将希望更有效地将其括起来(关闭现有的异常读者)。

【讨论】:

  • 有趣的文章,但它并没有真正谈论处理异常。您能否详细说明您将如何使用 scalaz Iteratees 来处理我的问题?
  • 我盯着这个看了 15 分钟,但我还是无法理解它。我想这意味着即使/当我弄清楚它时,编写这样的代码对我来说可能也不好......
  • 文章说明了。代码基本上说:要从 Reader 向 Iteratee 提供数据,请检查它是否已完成接受输入。如果是,请退回它。如果它期待更多的输入,它将有一个函数k 用于接受输入。从阅读器中读取一行并将其分配给s。如果我们得到一个异常,重新启动整个枚举。如果我们得到一个空行,则向 Iteratee 发出信号,告知我们已经到达 EOF。否则将s 提供给k 并循环。
  • 谢谢,这个解释比文章更有帮助。不过,我可能不会使用这种方法——它确实掩盖了我眼中的控制流。
  • 对,它对控制流进行了抽象。但是数据流很清晰。您正在以某种方式枚举字符串并生成 A。 “什么”是明确的,但“如何”是实现细节。
【解决方案3】:

这是一个不起作用的答案,但感觉应该:

def restartingIterator[T](getIter: () => Iterator[T]): Iterator[T] = {
  new Traversable[T] {
    def foreach[U](f: T => U): Unit = {
      try {
        for (item <- getIter()) {
          f(item)
        }
      } catch {
        case e: IOException => this.foreach(f)
      }
    }
  }.toIterator
}

我认为这非常清楚地描述了控制流,这很棒。

由于bug in Traversable.toStream,此代码将在Scala 2.8.0 中抛出StackOverflowError,但即使修复了该错误,此代码仍然不适用于我的用例,因为toIterator 调用@987654325 @,表示它会将所有项目存储在内存中。

我希望能够通过编写 foreach 方法来定义 Iterator,但似乎没有任何简单的方法可以做到这一点。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-03-14
    • 2013-02-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多