【问题标题】:Pyspark reduceByKey on nested tuples嵌套元组上的 Pyspark reduceByKey
【发布时间】:2016-12-27 12:38:22
【问题描述】:

我的问题与PySpark reduceByKey on multiple values 类似,但有一些关键的区别。我是 PySpark 的新手,所以我肯定会遗漏一些明显的东西。

我有一个具有以下结构的 RDD:

(K0, ((k01,v01), (k02,v02), ...))
....
(Kn, ((kn1,vn1), (kn2,vn2), ...))

我想要的输出类似于

(K0, v01+v02+...)
...
(Kn, vn1+vn2+...)

这似乎是使用reduceByKey 的完美案例,起初我想到了类似的东西

rdd.reduceByKey(lambda x,y: x[1]+y[1])

这正是我开始时的 RDD。我想我的索引有问题,因为有嵌套的元组,但我已经尝试了我能想到的所有可能的索引组合,并且它不断给我返回初始 RDD。

是否有它不应该与嵌套元组一起使用的原因,或者我做错了什么?

【问题讨论】:

    标签: python pyspark rdd reduce


    【解决方案1】:

    您根本不应该在这里使用reduceByKey。它需要一个带有签名的关联和交换函数。 (T, T) => T。很明显,当您将 List[Tuple[U, T]] 作为输入并且您希望 T 作为输出时,它不适用。

    由于不完全清楚键或唯一与否,让我们考虑当我们必须在本地和全局聚合时的一般示例。让我们假设 v01, v02, ... vm 是简单的数字:

    from functools import reduce
    from operator import add
    
    def agg_(xs):
        # For numeric values sum would be more idiomatic
        # but lets make it more generic
        return reduce(add, (x[1] for x in xs), zero_value)
    
    zero_value = 0
    merge_op = add
    def seq_op(acc, xs):
        return acc + agg_(xs)
    
    rdd = sc.parallelize([
        ("K0", (("k01", 3), ("k02", 2))),
        ("K0", (("k03", 5), ("k04", 6))),
        ("K1", (("k11", 0), ("k12", -1)))])
    
    rdd.aggregateByKey(0, seq_op, merge_op).take(2)
    ## [('K0', 16), ('K1', -1)]
    

    如果键已经是唯一的,简单的mapValues 就足够了:

    from itertools import chain
    
    unique_keys = rdd.groupByKey().mapValues(lambda x: tuple(chain(*x)))
    unique_keys.mapValues(agg_).take(2)
    ## [('K0', 16), ('K1', -1)]
    

    【讨论】:

    • 我现在很清楚了。是的,键是唯一的,所以 mapValues 方法正是我所需要的。非常感谢。
    猜你喜欢
    • 2014-01-31
    • 2015-07-02
    • 2015-10-17
    • 2020-03-14
    • 1970-01-01
    • 1970-01-01
    • 2018-08-07
    • 1970-01-01
    • 2021-10-26
    相关资源
    最近更新 更多