【问题标题】:A quick way to get the mean of each position in large RDD一种快速获取大型 RDD 中每个位置平均值的方法
【发布时间】:2023-04-09 11:17:01
【问题描述】:

我有一个很大的 RDD(超过 1,000,000 行),而每行有四个元素 A,B,C,D 在一个元组中。 RDD 的头部扫描看起来像

[(492,3440,4215,794),
(6507,6163,2196,1332),
(7561,124,8558,3975),
(423,1190,2619,9823)]

现在我想找到这个 RDD 中每个位置的平均值。例如,对于上面的数据,我需要一个输出列表具有值:

(492+6507+7561+423)/4
(3440+6163+124+1190)/4
(4215+2196+8558+2619)/4
(794+1332+3975+9823)/4

这是:

[(3745.75,2729.25,4397.0,3981.0)]

由于RDD很大,不方便计算每个位置的和然后除以RDD的长度。有什么快速的方法可以让我得到输出吗?非常感谢。

【问题讨论】:

    标签: python-3.x pyspark rdd


    【解决方案1】:

    我认为没有什么比计算每列的平均值(或总和)更快的方法了
    如果您使用的是 DataFrame API,您可以简单地聚合多个列:

    import os
    import time
    
    from pyspark.sql import functions as f
    from pyspark.sql import SparkSession
    
    # start local spark session
    spark = SparkSession.builder.getOrCreate()
    
    # load as rdd
    def localpath(path):
        return 'file://' + os.path.join(os.path.abspath(os.path.curdir), path)
    
    rdd = spark._sc.textFile(localpath('myPosts/'))
    
    # create data frame from rdd
    df = spark.createDataFrame(rdd)
    means_df = df.agg(*[f.avg(c) for c in df.columns])
    means_dict = means_df.first().asDict()
    print(means_dict)
    

    请注意,字典键将是默认的 spark 列名称('0'、'1'、...)。如果您想要更多说话的列名,可以将它们作为参数提供给 createDataFrame 命令

    【讨论】:

    • 您好,我的数据在 RDD 中。您能告诉我如何将其转换为 DataFrame 格式吗?非常感谢。
    • 我已经编辑并包含了如何从 RDD 创建 DF
    • 这里的spark 是什么?当我运行代码时出现错误:name 'spark' is not defined.
    • 这是 SparkSession:spark.apache.org/docs/2.1.0/api/python/… 如果你使用的是 spark 1.x,你需要一个 SqlContext
    • 能否请您写一个完整版的代码?由于我对 Spark 不熟悉,所以我也停留在构建会话部分。当我尝试使用spark = SparkSession.builder() 时出现另一个错误name 'SparkSession' is not defined
    猜你喜欢
    • 2015-08-03
    • 1970-01-01
    • 2019-01-08
    • 1970-01-01
    • 1970-01-01
    • 2017-09-07
    • 1970-01-01
    • 2012-12-30
    • 2022-11-10
    相关资源
    最近更新 更多