【问题标题】:Sum and divide elements of a RDD in pyspark在pyspark中对RDD的元素进行求和和除
【发布时间】:2021-11-29 21:06:33
【问题描述】:

我试图将 RDD 的所有元素相加,然后将其除以元素的数量。我能够解决它,但使用不同的行。但是我想只用一行使用 RDD 操作来完成。

RDD 例如:

rdd_example = [(eliana,1),(peter,2),(andrew,3),(paul,4),(jhon,5)]

第一步是使用带有 lambda 的方法映射来仅提取数字:

numbers = rdd_example.map(lambda x: x[1])

输出是:

numbers = [1,2,3,4,5]

然后是所有元素的总和,使用方法reduce:

from operator import add
sum = numbers.reduce(add)

然后使用 count 方法创建另一个变量来计算元素:

number_elem = rdd_example.count()

然后进行除法得到结果:

result = sum/number_elem 

我想只用一行和一个变量来完成所有这些工作。

【问题讨论】:

    标签: python apache-spark pyspark rdd


    【解决方案1】:

    使用fold,您可以一次性汇总计数和总和:

    cnt, total = rdd_example.fold((0, 0), lambda res, x: (res[0] + 1, res[1] + x[1]))
    
    print(total / cnt)
    # 2.5
    

    在调用中注意,我们使用一个元组来存储计数和总和:

    rdd_example.fold((0, 0), lambda res, x: (res[0] + 1, res[1] + x[1]))
    #                 ^  ^                   ^^^^^^^^^^  ^^^^^^^^^^^^^^
    #                 ^  init sum        add 1 to count / add value to sum
    #                 init count
    

    【讨论】:

    • 谢谢 Psidom 我真的很喜欢你的回答,方法 fold 是如何工作的?
    • 你可以在这里查看更多关于折叠的解释:spark.apache.org/docs/3.1.1/api/python/reference/api/…。基本上fold 需要一个初始值和一个 lambda。 lambda 从 rdd 中获取初始值和一个元素,计算并返回一个结果,该结果将成为新的初始值(或聚合值)。
    【解决方案2】:

    对于单线解决方案,请注意您正在计算数字的平均值(平均值)。 PySpark 已经有一个mean() 方法:

    rdd_example = sc.parallelize([("eliana",1),("peter",2),("andrew",3),("paul",4),("jhon",5)])
    result = rdd_example.map(lambda x: x[1]).mean()
    print(result)
    # output: 3.0
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-08-23
      • 1970-01-01
      • 2021-10-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多