【问题标题】:How to count the number of words per line in text file using RDD?如何使用 RDD 计算文本文件中每行的字数?
【发布时间】:2017-10-15 02:59:13
【问题描述】:

有没有办法使用 map 和 reduce 计算 RDD 的每一行的单词出现次数,而不是完整的 RDD?

例如,如果一个 RDD[String] 包含这两行:

让我们玩得开心。

为了玩得开心,你不需要任何计划。

那么输出应该像一个包含键值对的映射:

("让我们",1)
("有",1)
("一些",1)
(“有趣”,1)

("To",1)
("have",1)
("fun",1)
("you",1)
("don't", 1)
(“需要”,1)
(“计划”,1)

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您想要的是将一条线转换为 Map(word, count)。所以你可以逐行定义一个函数计数:

    def wordsCount(line: String):Map[String,Int] = {
     line.split(" ").map(v => (v,1)).groupBy(_._1).mapValues(_.size)
    }
    

    然后将其应用到您的 RDD[String]:

    val lines:RDD[String] = ...
    val wordsByLineRDD:RDD[Map[String,Int]] = lines.map(wordsCount)
    // this should give you a Map per line with count of each word
    wordsByLineRDD.take(2)
    // Something like
    // Array(Map(some -> 1, have -> 1, Let's -> 1, fun. -> 1), Map(any -> 1, have -> 1, don't -> 1, you -> 1, need -> 1, fun -> 1, To -> 1, plans. -> 1))
    

    【讨论】:

    • 能否请您为上述问题输入 pyspark 等效代码,即每行的字数
    • @Nandu查看最后一个帖子,有人提供了pyspark解决方案
    【解决方案2】:

    据我了解,您可以执行以下操作
    你说你有RDD[String]数据

    val data = Seq("Let's have some fun.",
      "To have fun you don't need any plans.")
    val rddData = sparkContext.parallelize(data)
    

    您可以将flatMap 应用于split lines 并在map 函数中创建(word, 1) tuples

    val output = rddData.flatMap(_.split(" ")).map(word => (word, 1))
    

    这应该会给你想要的输出

    output.foreach(println)
    

    要按行发生,您应该执行以下操作

    val output = rddData.map(_.split(" ").map((_, 1)).groupBy(_._1)
      .map { case (group: String, traversable) => traversable.reduce{(a,b) => (a._1, a._2 + b._2)} }.toList).flatMap(tuple => tuple)
    

    【讨论】:

    • 您的解决方案更简洁,但您没有按行计算出现次数
    • @KireetBhat,如果对您有帮助,请也点赞。谢谢
    • @RameshMaharjan 能否请您为每行的字数发布 pyspark 等效代码
    【解决方案3】:

    如果您刚开始使用 Spark,并且没有人告诉您使用它,请不要使用 RDD API。在 Spark 中,有很多更好且通常更高效的 Spark SQL API 可以执行此操作以及在大型数据集上执行许多其他分布式计算。

    使用 RDD API 就像将汇编程序用于可以使用 Scala(或其他高级编程语言)的东西。在开始你的 Spark 之旅时,我个人建议首先使用 DataFrames 和 Datasets 的更高级别的 Spark SQL API,这肯定是太多了。


    给定输入:

    $ cat input.txt
    Let's have some fun.
    
    To have fun you don't need any plans.
    

    如果您要使用 Dataset API,您可以执行以下操作:

    val lines = spark.read.text("input.txt").withColumnRenamed("value", "line")
    val wordsPerLine = lines.withColumn("words", explode(split($"line", "\\s+")))
    scala> wordsPerLine.show(false)
    +-------------------------------------+------+
    |line                                 |words |
    +-------------------------------------+------+
    |Let's have some fun.                 |Let's |
    |Let's have some fun.                 |have  |
    |Let's have some fun.                 |some  |
    |Let's have some fun.                 |fun.  |
    |                                     |      |
    |To have fun you don't need any plans.|To    |
    |To have fun you don't need any plans.|have  |
    |To have fun you don't need any plans.|fun   |
    |To have fun you don't need any plans.|you   |
    |To have fun you don't need any plans.|don't |
    |To have fun you don't need any plans.|need  |
    |To have fun you don't need any plans.|any   |
    |To have fun you don't need any plans.|plans.|
    +-------------------------------------+------+
    
    scala> wordsPerLine.
      groupBy("line", "words").
      count.
      withColumn("word_count", struct($"words", $"count")).
      select("line", "word_count").
      groupBy("line").
      agg(collect_set("word_count")).
      show(truncate = false)
    +-------------------------------------+------------------------------------------------------------------------------+
    |line                                 |collect_set(word_count)                                                       |
    +-------------------------------------+------------------------------------------------------------------------------+
    |To have fun you don't need any plans.|[[fun,1], [you,1], [don't,1], [have,1], [plans.,1], [any,1], [need,1], [To,1]]|
    |Let's have some fun.                 |[[have,1], [fun.,1], [Let's,1], [some,1]]                                     |
    |                                     |[[,1]]                                                                        |
    +-------------------------------------+------------------------------------------------------------------------------+
    

    完成。 很简单,不是吗?

    参见functions 对象(对于explodestruct 函数)。

    【讨论】:

      【解决方案4】:

      假设你有这样的rdd

      val data = Seq("Let's have some fun.",
        "To have fun you don't need any plans.")
      val rddData = sparkContext.parallelize(data)
      

      然后简单地申请flapMap 然后map

      val res = rddData.flatMap(line => line.split(" ")).map(word => (word,1))
      

      预期输出

      res.take(100)
      res4: Array[(String, Int)] = Array((Let's,1), (have,1), (some,1), (fun.,1), (To,1), (have,1), (fun,1), (you,1), (don't,1), (need,1), (any,1), (plans.,1))
      

      【讨论】:

        【解决方案5】:

        虽然这是一个老问题;我在 pySpark 中寻找答案。终于像下面这样管理了。

        file_ = cont_.parallelize (
            ["shots are shots that are shots with more big shots by big people",
             "people comes in all shapes and sizes, as people are idoits of the idiots",
             "i know what i am writing is nonsense, but i don't care because i am doing this to test my spark program",
             "my spark is a current spark, that spark in my eyes."]
        )
        
        file_ \
        .map(lambda x : [((x, i), 1) for i in x.split()]) \
        .flatMap(lambda x : x) \
        .reduceByKey(lambda x, y : x + y) \
        .sortByKey(False) \
        .map(lambda x : (x[0][1], x[1])) \
        .collect()
        

        【讨论】:

        • 这个问题被标记为scala,所以pySpark 解决方案不合适。请删除它。
        • 我遇到了这个问题,在搜索时这是我能找到的唯一相关页面。我看到它被标记为 Scala,但在主要问题中它是这样提到的。我提交了答案,以便其他可能正在搜索类似解决方案的人可能会发现它很有用。如果你觉得答案有误,请告诉我,我会改正的。如果你还是觉得完全没必要,那我就删了!
        猜你喜欢
        • 2017-02-16
        • 1970-01-01
        • 1970-01-01
        • 2017-01-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-07-18
        相关资源
        最近更新 更多