【问题标题】:Spark Scala: GroupByKey and sortSpark Scala:GroupByKey 和排序
【发布时间】:2016-04-30 03:07:46
【问题描述】:

我有一个RDD,其结构如下:

val rdd = RDD[ (category: String, product: String, score: Double) ]

我的目标是 group 基于类别的数据,然后为每个类别 sort w.r.t. Tuple 2 (product, score) 的分数。至于现在我的代码是:

val result = rdd.groupByKey.mapValues(v => v.toList.sortBy(-_._2))

事实证明,对于我拥有的数据而言,这是非常昂贵的操作。我希望使用替代方法来提高性能。

【问题讨论】:

  • 为什么排序如此重要?
  • 如果你能给出粗略的尺寸可能会有所帮助 - 原始 RDD 中有多少项目,有多少类别,平均每个类别有多少项目。这需要多长时间,在什么样的硬件上?您需要多快?
  • 您打算如何使用排序后的数据?你打算遍历所有这些,你只想找到最上面的吗?

标签: scala sorting apache-spark combiners


【解决方案1】:

在不知道您的数据集的情况下很难回答,但documentation 有一些线索回复:groupByKey 性能:

注意:此操作可能非常昂贵。如果你在分组 为了对每个执行聚合(例如总和或平均值) 键,使用 PairRDDFunctions.aggregateByKey 或 PairRDDFunctions.reduceByKey 将提供更好的性能。

所以这取决于您打算如何处理已排序的列表。如果您需要每个列表的全部内容,那么在groupByKey 上可能很难改进。如果您正在执行某种聚合,那么上面的替代操作(aggregateByKeyreduceByKey)可能会更好。

根据列表的大小,可能在排序之前使用替代集合(例如可变数组)会更有效。

编辑:如果你的类别比较少,你可以尝试反复过滤原始RDD,对每个过滤后的RDD进行排序。尽管总体上完成了类似的工作量,但它在任何特定时刻都可能使用更少的内存。

编辑 2:如果内存不足是个问题,您也许可以将您的类别和产品表示为整数 ID 而不是字符串,然后再查找名称。这样,您的主要 RDD 可能会小得多。

【讨论】:

  • 是的,我需要保留整个列表。这对应于一个商业案例,对于每个类别,我需要根据它们的排名列出产品。
【解决方案2】:

您的 RDD 在类别上是否公平分布?根据您的偏斜因素,您可能会遇到问题。 如果您没有太多键值,请尝试这样的操作:

val rdd: RDD[(String, String, Double)] = sc.parallelize(Seq(("someCategory","a",1.0),("someCategory","b",3.0),("someCategory2","c",4.0)))

rdd.keyBy(_._1).countByKey().foreach(println)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-04-23
    • 1970-01-01
    • 2023-03-26
    • 1970-01-01
    • 2016-05-07
    • 1970-01-01
    • 2017-06-13
    相关资源
    最近更新 更多