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 将底漆元素放到scan。 splitwhen 为每个起始范围开始一个新的子流。 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))