【问题标题】:Remove duplicate keys from Spark Scala从 Spark Scala 中删除重复键
【发布时间】:2015-07-28 16:21:09
【问题描述】:

我正在使用带有 scala 的 spark 1.2 并且有一对带有 (String, String) 的 RDD。示例记录如下所示:

<Key,  value>
id_1,  val_1_1; val_1_2
id_2,  val_2_1; val_2_2
id_3,  val_3_1; val_3_2
id_1,  val_4_1; val_4_2

我只是想删除所有有重复键的记录,所以在上面的例子中,第四条记录将被删除,因为 id_1 是重复键。

请帮忙。

谢谢。

【问题讨论】:

  • 哪里有重复键,你将如何决定保留哪个值?
  • 这只是我需要的第一个值。
  • 问题是,当 Spark 执行 reduceByKey 时,如下面的答案中所建议的,您无法知道将选择哪个值。不能保证 Spark 保持行的顺序。值(例如它是_1_1)有什么可以用来区分的吗?
  • @mattinbits,首先是 ziipWithIndex,然后在 reduce 中,只保留索引最低的那个,然后再映射以删除索引。中提琴!不,等等,那是一把大小提琴。沃利亚!

标签: scala apache-spark


【解决方案1】:

你可以使用reduceByKey:

val rdd: RDD[(K, V)] = // ...
val res: RDD[(K, V)] = rdd.reduceByKey((v1, v2) => v1)

【讨论】:

  • 谢谢。为什么它没有出现在我的脑海中。那很简单。谢谢。
【解决方案2】:

如果需要始终为给定键选择第一个条目,则将@JeanLogeart 答案与@Paul 的评论结合起来,

import org.apache.spark.{SparkContext, SparkConf}

val data = List(
  ("id_1", "val_1_1; val_1_2"),
  ("id_2",  "val_2_1; val_2_2"),
  ("id_3",  "val_3_1; val_3_2"),
  ("id_1",  "val_4_1; val_4_2") )

val conf = new SparkConf().setMaster("local").setAppName("App")
val sc = new SparkContext(conf)
val dataRDD = sc.parallelize(data)
val resultRDD = dataRDD.zipWithIndex.map{
  case ((key, value), index) => (key, (value, index))
}.reduceByKey((v1,v2) => if(v1._2 < v2._2) v1 else v2).mapValues(_._1)
resultRDD.collect().foreach(v => println(v))
sc.stop()

【讨论】:

    猜你喜欢
    • 2018-04-25
    • 1970-01-01
    • 2016-05-31
    • 1970-01-01
    • 1970-01-01
    • 2020-04-02
    • 1970-01-01
    • 1970-01-01
    • 2019-09-30
    相关资源
    最近更新 更多