【问题标题】:Scala Iterator/Looping Technique - Large CollectionsScala 迭代器/循环技术 - 大型集合
【发布时间】:2018-03-17 16:04:52
【问题描述】:

我有非常大的制表符分隔文件 (10GB-70GB),需要进行一些读取、数据操作和写入单独的文件。这些文件的范围可以从 100 到 10K 列,包含 200 万到 500 万行。

前 x 列是静态的,需要参考。示例文件格式:

#ProductName  Brand    Customer1  Customer2  Customer3
Corolla       Toyota   Y          N          Y
Accord        Honda    Y          Y          N
Civic         Honda    0          1          1

我需要使用前 2 列来获取产品 ID,然后生成类似于以下内容的输出文件:

ProductID1 Customer1 Y
ProductID1 Customer2 N
ProductID1 Customer3 Y
ProductID2 Customer1 Y
ProductID2 Customer2 Y
ProductID2 Customer3 N
ProductID3 Customer1 N
ProductID3 Customer2 Y
ProductID3 Customer3 Y

当前示例代码:

val fileNameAbsPath = filePath + fileName
val outputFile = new PrintWriter(filePath+outputFileName)
var customerList = Array[String]()

for(line <- scala.io.Source.fromFile(fileNameAbsPath).getLines()) {
    if(line.startsWith("#")) {
        customerList = line.split("\t")
    }
    else {
        val cols = line.split("\t")

        val productid = getProductID(cols(0), cols(1))
        for (i <- (2 until cols.length)) {
            val rowOutput = productid + "\t" + customerList(i) + "\t" + parser(cols(i))

            outputFile.println(rowOutput)
            outputFile.flush()
        }
    }
}
outputFile.close()

我运行的一个测试花了大约 12 个小时来读取一个包含 300 万行和 2500 列的文件 (70GB)。最终的输出文件生成了 250GB,大约有 800+ 百万行。

我的问题是:除了我已经在做的事情之外,Scala 中还有什么可以提供更快的性能吗?

【问题讨论】:

  • 这看起来很像一个作业,不适合作为一个问题
  • 如果我误解了,我很抱歉,但我只是在寻找想法,而不是有人帮助我编码。我对 scala 相当陌生,想知道它是否会提供更好的性能。
  • 我想将处理标题行的if 子句移出for 循环。如果您知道标题行只会出现一次,则无需对每一行执行检查。其次,除非你真的想确保你不想错过任何写入,flush 每次写入都会降低性能,我会先使用BufferedWriter 和@987654328 @ 并让他们负责刷新脏位。

标签: scala iterator


【解决方案1】:

好的,一些想法......

  • 如 cmets 中所述,您不想在每一行之后都添加flush。所以,是的,摆脱它。
  • 此外,PrintWriter 在默认情况下无论如何都会在每个换行符之后刷新(因此,目前,您实际上是在刷新两次 :))。在创建PrintWriter时使用双参数构造函数,并确保第二个参数为false
  • 您不需要显式创建BufferedWriterPrintWriter 默认情况下已经在缓冲。默认缓冲区大小为 8K,您可能想尝试使用它,但它可能不会有任何区别,因为,最后我检查,底层 FileOutputStream 忽略了所有这些,并以任何方式刷新千字节大小的块.
  • 无需在变量中将行粘合在一起,只需将每个字段直接写入输出即可。
  • 如果您不关心行在输出中出现的顺序,您可以简单地并行化处理(如果您确实关心顺序,您仍然可以,只是不那么琐碎一点),并编写几个文件立刻。如果您将输出块放在不同的磁盘上和/或如果您有多个内核来运行此代码,那将非常有帮助。您需要在(真正的)scala 中重写代码以使其线程安全,但这应该很容易。
  • 在写入数据时压缩数据。例如,使用GZipOutputStream。这不仅可以减少实际访问磁盘的物理数据量,还可以提供更大的缓冲区
  • 查看您的parser 正在做什么。您没有展示实现,但有消息告诉我它可能不是免费的。
  • split 在大字符串上可能会变得非常昂贵。人们经常忘记,它的参数实际上是一个正则表达式。您最好编写一个自定义迭代器或仅使用旧的 StringTokenizer 来解析字段,而不是预先拆分。至少,它会为您每行节省一次额外的扫描。

最后,最后,但绝不是最不重要的。 考虑使用 spark 和 hdfs。这类问题正是这些工具真正擅长的领域。

【讨论】:

  • 根据OpenJDK 8 sources,在单个字符上拆分相当快。
  • @OlegPyzhcov 是的......但是,您仍然必须创建一个包含 2500 列的数组,并预先填充它。如果 that 是你拒绝我的答案的原因,那么......你就是个混蛋 :)
  • 不,只是指出了一件奇怪且不太广为人知的事情(您的回答很好,我支持投票者:)
  • @OlegPyzhcov 好吧,很抱歉怀疑你 :)
猜你喜欢
  • 2015-02-08
  • 1970-01-01
  • 1970-01-01
  • 2014-01-11
  • 2018-11-18
  • 2016-11-19
  • 2011-02-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多