我相信climbage和eliasah忽略了问题的其他方面:
如果操作不减少数据量,它必须以一种或另一种语义等效于GroupByKey。假设我们有RDD[(Int,String)]:
import scala.util.Random
Random.setSeed(1)
def randomString = Random.alphanumeric.take(Random.nextInt(10)).mkString("")
val rdd = sc.parallelize((1 to 20).map(_ => (Random.nextInt(5), randomString)))
我们想要连接给定键的所有字符串。使用groupByKey 非常简单:
rdd.groupByKey.mapValues(_.mkString(""))
reduceByKey 的简单解决方案如下所示:
rdd.reduceByKey(_ + _)
它很短,可以说很容易理解,但存在两个问题:
- 效率极低,因为它每次都会创建一个新的
String 对象*
- 建议您执行的操作比实际操作更便宜,尤其是在您只分析 DAG 或调试字符串时
为了处理第一个问题,我们需要一个可变数据结构:
import scala.collection.mutable.StringBuilder
rdd.combineByKey[StringBuilder](
(s: String) => new StringBuilder(s),
(sb: StringBuilder, s: String) => sb ++= s,
(sb1: StringBuilder, sb2: StringBuilder) => sb1.append(sb2)
).mapValues(_.toString)
它仍然暗示了其他一些正在发生的事情并且非常冗长,特别是如果在您的脚本中重复多次。你当然可以提取匿名函数
val createStringCombiner = (s: String) => new StringBuilder(s)
val mergeStringValue = (sb: StringBuilder, s: String) => sb ++= s
val mergeStringCombiners = (sb1: StringBuilder, sb2: StringBuilder) =>
sb1.append(sb2)
rdd.combineByKey(createStringCombiner, mergeStringValue, mergeStringCombiners)
但归根结底,这仍然意味着要付出额外的努力来理解这段代码,增加复杂性并且没有真正的附加值。我发现特别令人不安的一件事是明确包含可变数据结构。即使 Spark 处理了几乎所有的复杂性,这也意味着我们不再拥有优雅的、引用透明的代码。
我的观点是,如果您真的想尽一切办法减少数据量,请使用reduceByKey。否则,您的代码会更难编写、更难分析,而且一无所获。
注意:
此答案主要针对 Scala RDD API。当前的 Python 实现与其对应的 JVM 有很大不同,并且包括一些优化,在类似 groupBy 的操作的情况下,这些优化比简单的 reduceByKey 实现具有显着优势。
对于Dataset API,请参阅DataFrame / Dataset groupBy behaviour/optimization。
* 请参阅Spark performance for Scala vs Python 以获得令人信服的示例