【问题标题】:Scala map-filtering methodsScala映射过滤方法
【发布时间】:2018-03-26 06:36:30
【问题描述】:

我是 Scala 和 Spark 的新手。我正在尝试删除文本文件的重复行。 每行包含三列(向量值),例如:-4.5,-4.2,2.7

这是我的程序:

import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.spark.rdd.RDD
import scala.collection.mutable.Map

object WordCount {

 def main(args: Array[String]) {

   val conf = new SparkConf().setAppName("WordCount").setMaster("local[*]")
   val sc = new SparkContext(conf)
   val input =  sc.textFile("/opt/spark/WC/WC_input.txt")

   val keys = input.flatMap(line => line.split("/n"))

   val singleKeys = keys.distinct

   singleKeys.foreach(println)
 }
}

它有效,但我想知道是否有一种方法可以使用过滤器功能。我必须在我的程序中使用它,但我不知道如何在所有行之间进行迭代并删除重复项(例如使用循环)。

如果有人有想法,那就太好了!

谢谢!

【问题讨论】:

    标签: scala apache-spark filter rdd


    【解决方案1】:

    我认为使用filter 这样做不是一个非常有效的解决方案。对于每个元素,您必须要么查看该元素是否已经存在于某种临时数据集中,要么计算这些元素中有多少在已处理的数据集中。

    如果您想对其进行迭代并可能进行一些即时编辑,您可以应用map 然后reduceByKey 对相同的元素进行分组。像这样

    val singleKeys = 
        keys
        .map( element => ( element , 0 ) )
        .reduceByKey( ( element, count ) => element )
        .map( _._1 )
    

    您可以在第一个map 部分中对数据集进行更改。 count 参数未使用,尽管根据 reduceByKey 的定义,我们需要 Tuple 或 Map 中的第二个参数。

    我认为这基本上就是 distinct 内部的工作方式。

    【讨论】:

      【解决方案2】:

      RDD的重复元素可以这样去除:

      val data = List("-4.5,-4.2,2.7", "10,20,30", "-4.5,-4.2,2.7")
      val rdd = sparkContext.parallelize(data)
      val result = rdd.map((_, 1)).reduceByKey(_ + _).filter(_._2 == 1).map(_._1)
      result.foreach(println)
      

      结果:

      10,20,30
      

      【讨论】:

      • 非常感谢!!但我想保留所有重复元素的一个实例。有可能吗?
      • 是的,如果删除“过滤器”子句。结果将与“distinct”相同。
      • 你是最棒的!谢谢(:
      猜你喜欢
      • 2015-08-30
      • 1970-01-01
      • 2022-07-06
      • 2016-07-22
      • 1970-01-01
      • 1970-01-01
      • 2015-11-20
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多