【问题标题】:Reduce job in Spark by reduceByKey() or other functions?通过 reduceByKey() 或其他函数减少 Spark 中的作业?
【发布时间】:2015-12-28 02:43:39
【问题描述】:

给定以下列表:

[(0, [135, 2]), (0, [2409, 1]), (0, [12846, 2]), (1, [13840, 2]), ...]

如果列表值的第二个元素,我需要为每个键输出列表值的第一个元素的列表(即,135, 2409, 12846 用于键 0 和 13840 用于键 1) (即,2, 1, 2 用于0 和2 用于1)大于或等于某个值(假设为2)。例如,在这种特殊情况下,输出应该是:

[(0, [135, 12846]), (1, [13840]), ...]

元组(0, [2409, 1]) 被丢弃,因为1 < 2。

我通过应用groupByKey()、mapValues(list) 和最终的map 函数实现了这一点,但显然groupByKey() 的效率低于reduce 函数。

仅使用reduceByKey() 或combineByKey() 函数就可以完成该任务吗?

【问题讨论】:

    标签: python mapreduce apache-spark reduce pyspark


    【解决方案1】:

    答案是肯定的 :) 您可以使用 reduceByKey 和 groupByKey 实现相同的效果。事实上,reduceByKey 应该始终受到青睐,因为它在洗牌数据之前会执行 map 侧归约。

    使用reduceByKey 的解决方案(在 Scala 中,但我相信您明白这一点,如果您愿意,可以轻松地将其转换为 Python):

    val rdd = sc.parallelize(List((0, List(135, 2)), (0, List(2409, 1)), (0, List(12846, 2)), (1, List(13840, 2))))
    rdd.mapValues(v => if(v(1) >= 2) List(v(0)) else List.empty)
       .reduceByKey(_++_)
    

    【讨论】:

    • 谢谢格伦尼。只是一个更好地理解如何在 Python 代码中翻译 Scala 代码的问题:++ 在 scala 中是什么意思?这是一个简单的总和吗?
    • ++ 是一个列表连接。在 Python 中只需使用 +。
    • Glennie,always 在这里用词太强了。当函数可以减少数据量时,像 这样的东西是首选 会是更好的选择。
    • @zero323 是的,你是对的。我想我只是想鼓励人们在选择groupByKey 之前考虑reduceByKey。我发现groupByKey 通常被使用的太多了。
    • 完美,再次感谢!我已将 scala 代码翻译成 python 代码:rdd.mapValues(lambda x: [x[0]] if x[1] >= 3 else []).reduceByKey(lambda a, b: a + b)
    猜你喜欢
    • 2016-02-25
    • 2017-04-10
    • 2023-01-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-02-07
    • 2017-04-03
    • 2021-09-06
    相关资源
    最近更新 更多