【问题标题】:Spark reduceByKey on several different valuesSpark reduceByKey 在几个不同的值上
【发布时间】:2015-07-06 22:36:32
【问题描述】:

我有一个存储为 RDD 列表的表,我想在其上执行类似于 SQL 或 pandas 中的 groupby 的操作,获取每个变量的总和或平均值。

我目前的做法是这样的(未经测试的代码):

l=[(3, "add"),(4, "add")]
dict={}
i=0
for aggregation in l:
    RDD= RDD.map(lambda x: (x[6], float(x[aggregation[0]])))
    agg=RDD.reduceByKey(aggregation[1])
    dict[i]=agg
    i+=1

然后我需要加入dict中的所有RDD。

虽然这不是很有效。有没有更好的办法?

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    如果您使用 >= Spark 1.3,您可以查看DataFrame API

    在 pyspark 外壳中:

    import numpy as np
    # create a DataFrame (this can also be from an RDD)
    df = sqlCtx.createDataFrame(map(lambda x:map(float, x), np.random.rand(50, 3)))
    df.agg({col: "mean" for col in df.columns}).collect()
    

    这个输出:

    [Row(AVG(_3#1456)=0.5547187588389414, AVG(_1#1454)=0.5149476209374797, AVG(_2#1455)=0.5022967093047612)]
    

    可用的聚合方法有“avg”/“mean”、“max”、“min”、“sum”、“count”。

    要为同一列获取多个聚合,您可以使用显式构造的聚合列表而不是字典调用agg

    from pyspark.sql import functions as F
    df.agg(*[F.min(col) for col in df.columns] + [F.avg(col) for col in df.columns]).collect()
    

    或者对于您的情况:

    df.agg(F.count(df.var3), F.max(df.var3), ) # etc...
    

    【讨论】:

    • 您能否启发我对特定变量的特定聚合方法的语法。目前下面只给了我对每个变量使用的最后一种方法: agggreg=df1.filter(df1.var3>25).join(df2, df1.Machine== df2.Machine).groupBy(df1.Machine).agg ({"var3":"count","var3":"max","var4":"mean","var4":"max","var5":"mean","var6":"mean"} )
    • 谢谢。使用这种新语法,groupby 变量(在本例中为机器)消失了。我怎样才能保留它?
    猜你喜欢
    • 2015-05-04
    • 2020-02-08
    • 2017-10-07
    • 2018-01-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-30
    • 2014-07-19
    相关资源
    最近更新 更多