【问题标题】:Sink for line-by-line file IO with backpressure带背压的逐行文件 IO 接收器
【发布时间】:2016-04-15 21:06:35
【问题描述】:

我有一个文件处理作业,目前使用具有手动管理背压的 akka Actor 来处理处理管道,但我从未能够在输入文件读取阶段成功管理背压。

此作业获取一个输入文件并根据每行开头的 ID 号对行进行分组,然后一旦遇到具有新 ID 号的行,它就会通过消息将分组的行推送给处理参与者,并且然后继续使用新的 ID 号,一直到文件末尾。

这似乎是 Akka Streams 的一个很好的用例,使用 File 作为接收器,但我仍然不确定三件事:

1) 如何逐行读取文件?

2) 如何按每行的 ID 进行分组?我目前为此使用非常命令式的处理,我认为我不会在流管道中拥有相同的能力。

3) 我怎样才能应用背压,这样我就不会比我处理下游数据的速度更快地将行读入内存?

【问题讨论】:

  • 问题:您现在如何管理背压?你是从单个节点读取文件吗?你是用akka集群来处理的吗?
  • 我不管理背压。我已经尝试了一些东西,但它们都很老套(比如长轮询asking 来自处理参与者的Continue 消息,并手动迭代读取,这非常脆弱)。所以我选择通过给我的应用程序足够的堆空间来将整个输入文件消耗到内存中来破解它。我不能再这样做了,因为我必须将它部署到共享服务器上,而且我不能再吞噬每个人的内存了。

标签: scala akka akka-stream


【解决方案1】:

Akka 流的groupBy 是一种方法。但是 groupBy 有一个 maxSubstreams 参数,它要求您预先知道最大 ID 范围。所以:下面的解决方案使用scan标识同ID块,splitWhen拆分成子流:

object Main extends App {
  implicit val system = ActorSystem("system")
  implicit val materializer = ActorMaterializer()

  def extractId(s: String) = {
    val a = s.split(",")
    a(0) -> a(1)
  }

  val file = new File("/tmp/example.csv")

  private val lineByLineSource = FileIO.fromFile(file)
    .via(Framing.delimiter(ByteString("\n"), maximumFrameLength = 1024))
    .map(_.utf8String)

  val future: Future[Done] = lineByLineSource
    .map(extractId)
    .scan( (false,"","") )( (l,r) => (l._2 != r._1, r._1, r._2) )
    .drop(1)
    .splitWhen(_._1)
    .fold( ("",Seq[String]()) )( (l,r) => (r._2, l._2 ++ Seq(r._3) ))
    .concatSubstreams
    .runForeach(println)

  private val reply = Await.result(future, 10 seconds)
  println(s"Received $reply")
  Await.ready(system.terminate(), 10 seconds)
}

extractId 将行拆分为 id -> 数据元组。 scan 在 id -> data tuples 前面加上一个 start-of-ID-range 标志。 drop 将底漆元素放到scansplitwhen 为每个起始范围开始一个新的子流。 fold 将子流连接到列表并删除起始 ID 范围布尔值,以便每个子流生成单个元素。您可能需要一个自定义的 SubFlow 来代替折叠,它处理单个 ID 的行流并为 ID 范围发出一些结果。 concatSubstreams 将 splitWhen 生成的 per-ID-range 子流合并回 runForEach 打印的单个流。

运行:

$ cat /tmp/example.csv
ID1,some input
ID1,some more input
ID1,last of ID1
ID2,one line of ID2
ID3,2nd before eof
ID3,eof

输出是:

(ID1,List(some input, some more input, last of ID1))
(ID2,List(one line of ID2))
(ID3,List(2nd before eof, eof))

【讨论】:

  • 我想我在这里遵循逻辑,但你能解释一下mergeSubstreams 的作用吗?我似乎无法在 API 文档中找到它的定义。
  • 我没有找到除了 scaladocs 之外的任何文档,但我为舞台添加了 cmets。另外,我用concatSubstreams 替换了mergeSubstreams,因为在这种情况下,子流按顺序启动和终止(根据ID 范围),并且无论子流中的任何异步处理如何,concat 都会以与ID 范围相同的顺序产生输出。
【解决方案2】:

似乎在不引入大量修改的情况下向系统添加“背压”的最简单方法是将使用 Actor 的输入组的邮箱类型更改为 BoundedMailbox

  1. 将消耗您的台词的 Actor 的类型更改为 BoundedMailbox 并具有高 mailbox-push-timeout-time

    bounded-mailbox {
      mailbox-type = "akka.dispatch.BoundedDequeBasedMailbox"
      mailbox-capacity = 1
      mailbox-push-timeout-time = 1h
    }
    
    val actor = system.actorOf(Props(classOf[InputGroupsConsumingActor]).withMailbox("bounded-mailbox"))
    
  2. 从您的文件创建迭代器,从该迭代器创建分组(按 id)迭代器。然后只需循环遍历数据,将组发送到消费 Actor。请注意,在这种情况下,当 Actor 的邮箱已满时,发送将被阻塞。

    def iterGroupBy[A, K](iter: Iterator[A])(keyFun: A => K): Iterator[Seq[A]] = {
      def rec(s: Stream[A]): Stream[Seq[A]] =
        if (s.isEmpty) Stream.empty else {
          s.span(keyFun(s.head) == keyFun(_)) match {
          case (prefix, suffix) => prefix.toList #:: rec(suffix)
        }
      }
      rec(iter.toStream).toIterator
    }
    
    val lines = Source.fromFile("input.file").getLines()
    
    iterGroupBy(lines){l => l.headOption}.foreach {
      lines:Seq[String] =>
          actor.tell(lines, ActorRef.noSender)
    }
    

就是这样! 您可能希望将文件读取内容移动到单独的线程,因为它会阻塞。此外,通过调整mailbox-capacity,您可以调节消耗的内存量。但是,如果从文件中读取批处理总是比处理更快,那么保持较小的容量(例如 1 或 2)似乎是合理的。

更新 iterGroupByStream 一起实现,经测试不会产生StackOverflow

【讨论】:

    猜你喜欢
    • 2018-08-24
    • 2014-09-23
    • 2020-08-26
    • 2019-06-20
    • 2015-06-08
    • 1970-01-01
    • 1970-01-01
    • 2012-10-20
    • 1970-01-01
    相关资源
    最近更新 更多