【问题标题】:How can count words of multiple files present in a directory using spark scala如何使用 spark scala 计算目录中存在的多个文件的单词
【发布时间】:2018-11-28 21:02:24
【问题描述】:

如何使用 Apache Spark 和 Scala 对目录中存在的多个文件进行字数统计?

所有文件都有换行符。

O/p 应该是:

file1.txt,5
file2.txt,6 ...

我尝试使用以下方式:

val rdd= spark.sparkContext.wholeTextFiles("file:///C:/Datasets/DataFiles/")
val cnt=rdd.map(m =>( (m._1,m._2),1)).reduceByKey((a,b)=> a+b)

O/p 我得到了:

((file:/C:/Datasets/DataFiles/file1.txt,apple
orange
bag
apple
orange),1)
((file:/C:/Datasets/DataFiles/file2.txt,car
bike
truck
car
bike
truck),1)

我首先尝试了sc.textFile(),但没有给我文件名。 wholeTextFile() 返回键值对,其中键是文件名,但无法得到想要的输出。

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您的起步是正确的,但需要在您的解决方案上多做一些工作。

    sparkContext.wholeTextFiles(...) 方法为您提供了一个 (file, contents) 对,因此当您通过密钥减少它时,您会得到 (file, 1),因为这是每个密钥对拥有的整个文件内容的数量.

    为了统计每个文件的单词,你需要将每个文件的内容分解成那些单词,这样你就可以统计它们了。

    就到这里吧,我们开始读取文件目录:

    val files: RDD[(String, String)] = spark.sparkContext.wholeTextFiles("file:///C:/Datasets/DataFiles/")
    

    这为每个文件提供一行,以及完整的文件内容。现在让我们将文件内容分解为单个项目。鉴于您的文件似乎每行只有一个单词,因此使用换行符非常容易:

    val wordsPerFile: RDD[(String, Array[String])] = files.mapValues(_.split("\n"))
    

    现在我们只需要计算每个 Array[String] 中存在的项目数:

    val wordCountPerFile: RDD[(String, Int)] = wordsPerFile.mapValues(_.size)
    

    基本上就是这样。值得一提的是,虽然 字数统计 根本没有分发(它只是使用 Array[String]),因为您正在一次加载文件的全部内容。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-14
      • 2021-06-25
      • 2016-06-04
      • 2012-08-04
      相关资源
      最近更新 更多