【发布时间】:2018-12-25 03:07:11
【问题描述】:
我有这样的来源:
1 | red | light | 10
2 | blue | dark | 20
1 | brown | light | 2
1 | red | light | 10
20 | grey | dark | 200
我想知道 (true / false) 源中是否有任何相同的项目。在上面的流中1 | red | light | 10 将是相同的。这个流可能非常大,超过 2M 条记录。我可以在找到相同的项目后立即返回true(即在上面的示例中,我们可以避免阅读20 | grey | dark | 200)。
最好的方法是什么?我尝试将整个源代码读入List(String) 并在其上运行不同。这可以正常工作,但是,对于大型来源,我开始收到 OOM 错误。
val restResult: Future[immutable.Seq[Color]] =
mySource(ctx)
.drop(1)
.via(framing("\n"))
.map(_.utf8String)
.map(_.trim)
.map(s => ColorParser(s))
.collect {
case Right(color) => color
}
.runWith(Sink.seq)
【问题讨论】:
-
如何使用流逐行读取 CSV 并将元素递归地附加到累加器中,并在读取新元素时在累加器中进行检查?
-
我正在逐行读取元素(提供代码)。我不熟悉累加器的概念。你能详细说明一下吗?
-
不确定 Akka-streams 是否是正确的工具。如果您没有找到任何重复项,您可以轻松获得 StackOverflowException。 Spark 或 Flink 看起来是更好的选择。
标签: scala akka akka-stream