【问题标题】:Scala spark reduce by key and find common valueScala spark按键减少并找到共同价值
【发布时间】:2017-01-07 18:52:16
【问题描述】:

我有一个 csv 数据文件存储在 HDFS 上的 sequenceFile 中,格式为 name, zip, country, fav_food1, fav_food2, fav_food3, fav_colour。可能有许多具有相同名称的条目,我需要找出他们最喜欢的食物是什么(即计算所有具有该名称的记录中的所有食物条目并返回最受欢迎的条目。我是 Scala 和 Spark 的新手,并且有浏览了多个教程并搜索了论坛,但一直不知道如何继续。到目前为止,我已经得到了将文本转换为字符串格式的序列文件,然后过滤掉了条目

这是文件中每一行的示例数据条目

Bob,123,USA,Pizza,Soda,,Blue
Bob,456,UK,Chocolate,Cheese,Soda,Green
Bob,12,USA,Chocolate,Pizza,Soda,Yellow
Mary,68,USA,Chips,Pasta,Chocolate,Blue

所以输出应该是元组 (Bob, Soda),因为 soda 在 Bob 的条目中出现的次数最多。

import org.apache.hadoop.io._

var lines  = sc.sequenceFile("path",classOf[LongWritable],classOf[Text]).values.map(x => x.toString())
// converted to string since I could not get filter to run on Text and removing the longwritable

var filtered = lines.filter(_.split(",")(0) == "Bob");
// removed entries with all other users

var f_tuples = filtered.map(line => lines.split(",");
// split all the values

var f_simple = filtered.map(line => (line(0), (line(3), line(4), line(5))
// removed unnecessary fields

我现在遇到的这个问题是,我认为我有这个 [<name,[f,f,f]>] 结构,但我真的不知道如何将其展平并获得最受欢迎的食物。我需要合并所有条目,所以我有一个带有 a 的条目,然后获取值中最常见的元素。任何帮助,将不胜感激。谢谢

我试过这个让它变平,但似乎我尝试得越多,数据结构就越复杂。

var f_trial = fpairs.groupBy(_._1).mapValues(_.map(_._2))
// the resulting structure was of type org.apache.spark.rdd.RDD[(String, Interable[(String, String, String)]

这是 f_trial 之后记录的 println 的样子

("Bob", List((Pizza, Soda,), (Chocolate, Cheese, Soda), (Chocolate, Pizza, Soda)))

括号分解

("Bob", 

List(

(Pizza, Soda, <missing value>),

(Chocolate, Cheese, Soda),

(Chocolate, Pizza, Soda)

) // ends List paren

) // ends first paren

【问题讨论】:

  • 您是否需要为每个人提供一种最受欢迎​​的食物,或者您是否需要每个名字都获得最喜欢的食物?您能否提供一个示例 od 您的 f_simple 数据以及您想要得到什么?
  • @Niemand 现在我只需要该名称的名称和最受欢迎的食物(将来我可能需要获得最受欢迎的前 3 名或其他什么,但我只需要一个基础即可开始on),我也在想同样的事情,我应该提供数据集并且刚刚提供,谢谢回复。
  • Bob 的第一个条目只有两种最喜欢的食物。记录的列数是否都相同?
  • @Paul,感谢回复,每个条目最多可以有3种喜欢的食物,我可以少但绝对不能多。即使只有 1 或 2 个条目(只是在其中编辑),也会包含尾随逗号。
  • 没有时间得到正确的答案,但我认为你需要扁平化到(name, food)(或者可能是((name, food), 1))然后reducebyKey((name,food),total)Map(name, (food, total)) , reduceByKey 再次(仅在减少步骤中保留最大总数)。这为您提供(name, (food, total)) 每个人最受欢迎的食物

标签: scala hadoop apache-spark


【解决方案1】:

我找到了时间。设置:

    import org.apache.spark.SparkContext
    import org.apache.spark.SparkContext._
    import org.apache.spark.SparkConf

    val conf = new SparkConf().setAppName("spark-scratch").setMaster("local")
    val sc = new SparkContext(conf)

    val data = """   
  Bob,123,USA,Pizza,Soda,,Blue
  Bob,456,UK,Chocolate,Cheese,Soda,Green
  Bob,12,USA,Chocolate,Pizza,Soda,Yellow
  Mary,68,USA,Chips,Pasta,Chocolate,Blue
  """.trim

    val records = sc.parallelize(data.split('\n'))

提取食物选择,并为每个选择一个 ((name, food), 1) 元组

    val r2 = records.flatMap { r =>
      val Array(name, id, country, food1, food2, food3, color) = r.split(',');
      List(((name, food1), 1), ((name, food2), 1), ((name, food3), 1))
    }

合计每个名称/食物组合:

    val r3 = r2.reduceByKey((x, y) => x + y)

重新映射,以便名称(仅)是关键

    val r4 = r3.map { case ((name, food), total) => (name, (food, total)) }

在每一步中选择数量最多的食物

    val res = r4.reduceByKey((x, y) => if (y._2 > x._2) y else x)

我们完成了

    println(res.collect().mkString)
    //(Mary,(Chips,1))(Bob,(Soda,3))

编辑:要收集对一个人具有相同最高计数的所有食物,我们只需更改最后两行:

从包含总计的项目列表开始:

val r5 = r3.map { case ((name, food), total) => (name, (List(food), total)) }

在相同的情况下,将带有该分数的食物列表连接起来

val res2 = r5.reduceByKey((x, y) => if (y._2 > x._2) y 
                                    else if (y._2 < x._2) x
                                    else (y._1:::x._1, y._2))

//(Mary,(List(Chocolate, Pasta, Chips),1))
//(Bob,(List(Soda),3))

如果你想要前三名,那么使用aggregateByKey 来组合每个人最喜欢的食物列表,而不是第二个reduceByKey

【讨论】:

  • 这太棒了。我以错误的方式处理问题,我想从更传统的角度来看。你能解释一下当涉及到相同数量的食物时你会如何处理领带吗?现在只显示一个。
  • 当你有相同数量的食物时,你想要什么结果?最简单的方法是建立一个具有相同(最大)数量的食物列表。
  • 那将是完美的,我有另一个解析程序可以获取一个列表并稍后对其进行操作,其逻辑将在哪里?我在使用 aggregateByKey 时也遇到了问题。我非常感谢你们的所有帮助
  • 好的,完成了。不确定我是否理解您的第二点,但如果您遇到困难,不妨尝试自己并发布另一个问题?
  • 别忘了给有帮助的答案点赞! :-P
【解决方案2】:

Paulmattinbits 提供的解决方案将您的数据洗牌两次 - 一次执行 reduce-by-name-and-food,一次执行 reduce-by-name。只需一次 shuffle 即可解决此问题。

/**Generate key-food_count pairs from a splitted line**/
def bitsToKeyMapPair(xs: Array[String]): (String, Map[String, Long]) = {
  val key = xs(0)
  val map = xs
    .drop(3) // Drop name..country
    .take(3) // Take food
    .filter(_.trim.size !=0) // Ignore empty
    .map((_, 1L)) // Generate k-v pairs
    .toMap // Convert to Map
    .withDefaultValue(0L) // Set default

  (key, map)
}

/**Combine two count maps**/
def combine(m1: Map[String, Long], m2: Map[String, Long]): Map[String, Long] = {
  (m1.keys ++ m2.keys).map(k => (k, m1(k) + m2(k))).toMap.withDefaultValue(0L)
}

val n: Int = ??? // Number of favorite per user

val records = lines.map(line => bitsToKeyMapPair(line.split(",")))
records.reduceByKey(combine).mapValues(_.toSeq.sortBy(-_._2).take(n))

如果您不是纯粹主义者,可以将scala.collection.immutable.Map 替换为scala.collection.mutable.Map 以进一步提高性能。

【讨论】:

  • 效率更高,是的。但不是很清楚,我认为。就像@mattinbits,.toSeq.sortBy(-_._2).take(1) is maxBy(._2)。我的另一个优点是,如果有很多很多可供选择的食物,效果会更好(我承认,不太可能!)
  • sortByKey 允许您获取任意数量的元素,并且集合 API 不提供 topBy 方法。关于大量选项,我不能同意。假设用户有 N 个“最喜欢”的食物,每个食物在数据集中只出现一次。您必须传输此数据两次,唯一的好处是在计算最终答案时减少了内存占用。
  • 问题是如果一个用户的食物列表不适合一个节点(不太可能,正如我所说)。此外,它是maxBy,而不是toBy,并提供:scala-lang.org/api/current/…。您不需要“任意数量的元素”,您需要 1(因此您的 `take(1)`。
  • 我已经编辑了我的答案以明确我的意图:) 关于不适合一个节点,这是极不可能的。假设每个工作人员 16GB 和(每对 64 字节 * 对数)+ 映射对象标头,如果我计数正确,您可以存储 1e8 对,为组合留出足够的空间:)
  • 是的,不太可能。为了清楚起见,我仍然更喜欢我的,但是对于真实世界的数据集,你的效率很好。遗憾的是 Scala 集合 API 没有 top(而 Spark 没有 topByKey - 对于普通 RDD,它确实有`top`)
【解决方案3】:

这是一个完整的例子:

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


object Main extends App {

  val data = List(
    "Bob,123,USA,Pizza,Soda,,Blue",
    "Bob,456,UK,Chocolate,Cheese,Soda,Green",
    "Bob,12,USA,Chocolate,Pizza,Soda,Yellow",
    "Mary,68,USA,Chips,Pasta,Chocolate,Blue")

  val sparkConf = new SparkConf().setMaster("local").setAppName("example")
  val sc = new SparkContext(sparkConf)

  val lineRDD = sc.parallelize(data)

  val pairedRDD = lineRDD.map { line =>
    val fields = line.split(",")
    (fields(0), List(fields(3), fields(4), fields(5)).filter(_ != ""))
  }.filter(_._1 == "Bob")

  /*pairedRDD.collect().foreach(println)
    (Bob,List(Pizza, Soda))
    (Bob,List(Chocolate, Cheese, Soda))
    (Bob,List(Chocolate, Pizza, Soda))
   */

  val flatPairsRDD = pairedRDD.flatMap {
    case (name, foodList) => foodList.map(food => ((name, food), 1))
  }

  /*flatPairsRDD.collect().foreach(println)
    ((Bob,Pizza),1)
    ((Bob,Soda),1)
    ((Bob,Chocolate),1)
    ((Bob,Cheese),1)
    ((Bob,Soda),1)
    ((Bob,Chocolate),1)
    ((Bob,Pizza),1)
    ((Bob,Soda),1)
   */

  val nameFoodSumRDD = flatPairsRDD.reduceByKey((a, b) => a + b)

  /*nameFoodSumRDD.collect().foreach(println)
    ((Bob,Cheese),1)
    ((Bob,Soda),3)
    ((Bob,Pizza),2)
    ((Bob,Chocolate),2)
   */

  val resultsRDD = nameFoodSumRDD.map{
    case ((name, food), count) => (name, (food,count))
  }.groupByKey.map{
    case (name, foodCountList) => (name, foodCountList.toList.sortBy(_._2).reverse.head)
  }

  resultsRDD.collect().foreach(println)
  /*
      (Bob,(Soda,3))
   */

  sc.stop()
}

【讨论】:

  • 这有点相似,而不是我的方法的“完整示例”。注意我故意避开groupByKey。来自 Spark 文档:“注意:此操作可能非常昂贵。如果您正在分组以便对每个键执行聚合(例如总和或平均值),则使用 PairRDDFunctions.aggregateByKeyPairRDDFunctions.reduceByKey 将提供更好的性能。”
  • 而这个foodCountList.toList.sortBy(_._2).reverse.head 是一种复杂的表达方式maxBy (_._2)
  • 我引用了您的评论公司。 “没有时间给出正确的答案”不是你的答案,你的答案一定在我写我的时候上升了,道歉。
  • 没问题。很高兴给 OP 一些替代方案
  • @mattinbits 这与 Paul 的非常相似,尽管我最终使用这里的逻辑来过滤掉空的食物条目,并使用基于索引的访问,因为我在期末考试中有更多的字段项目。谢谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-09-17
  • 1970-01-01
  • 2016-12-05
  • 2018-11-11
  • 2016-11-24
  • 1970-01-01
  • 2018-10-09
相关资源
最近更新 更多