【问题标题】:How to convert groupby into reducebykey in pyspark dataframe? [duplicate]如何在pyspark数据框中将groupby转换为reducebykey? [复制]
【发布时间】:2018-03-02 01:08:21
【问题描述】:

我用 group by 和 sum 函数编写了 pyspark 代码。由于分组,我觉得性能受到影响。相反,我想使用 reducebykey。但我是这个领域的新手。请在下面找到我的方案,

第一步:通过sqlcontext读取hive表连接查询数据并存储在dataframe中

Step2:输入列的总数为15。其中5个是关键字段,其余是数值。

第 3 步:除了上述输入列之外,还需要从数字列中派生出更多列。很少有具有默认值的列。

第 4 步:我使用了 group by 和 sum 函数。如何使用带有 map 和 reducebykey 选项的 spark 方式执行类似的逻辑。

from pyspark.sql.functions import col, when, lit, concat, round, sum

#sample data
df = sc.parallelize([(1, 2, 3, 4), (5, 6, 7, 8)]).toDF(["col1", "col2", "col3", "col4"])

#populate col5, col6, col7
col5 = when((col('col1') == 0) & (col('col3') != 0), round(col('col4')/ col('col3'), 2)).otherwise(0)
col6 = when((col('col1') == 0) & (col('col4') != 0), round((col('col3') * col('col4'))/ col('col1'), 2)).otherwise(0)
col7 = col('col2')
df1 = df.withColumn("col5", col5).\
    withColumn("col6", col6).\
    withColumn("col7", col7)

#populate col8, col9, col10
col8 = when((col('col1') != 0) & (col('col3') != 0), round(col('col4')/ col('col3'), 2)).otherwise(0)
col9 = when((col('col1') != 0) & (col('col4') != 0), round((col('col3') * col('col4'))/ col('col1'), 2)).otherwise(0)
col10= concat(col('col2'), lit("_NEW"))
df2 = df.withColumn("col5", col8).\
    withColumn("col6", col9).\
    withColumn("col7", col10)

#final dataframe
final_df = df1.union(df2)
final_df.show()

#groupBy calculation
#final_df.groupBy("col1", "col2", "col3", "col4").agg(sum("col5")).show()from pyspark.sql.functions import col, when, lit, concat, round, sum

#sample data
df = sc.parallelize([(1, 2, 3, 4), (5, 6, 7, 8)]).toDF(["col1", "col2", "col3", "col4"])

#populate col5, col6, col7
col5 = when((col('col1') == 0) & (col('col3') != 0), round(col('col4')/ col('col3'), 2)).otherwise(0)
col6 = when((col('col1') == 0) & (col('col4') != 0), round((col('col3') * col('col4'))/ col('col1'), 2)).otherwise(0)
col7 = col('col2')
df1 = df.withColumn("col5", col5).\
    withColumn("col6", col6).\
    withColumn("col7", col7)

#populate col8, col9, col10
col8 = when((col('col1') != 0) & (col('col3') != 0), round(col('col4')/ col('col3'), 2)).otherwise(0)
col9 = when((col('col1') != 0) & (col('col4') != 0), round((col('col3') * col('col4'))/ col('col1'), 2)).otherwise(0)
col10= concat(col('col2'), lit("_NEW"))
df2 = df.withColumn("col5", col8).\
    withColumn("col6", col9).\
    withColumn("col7", col10)

#final dataframe
final_df = df1.union(df2)
final_df.show()

#groupBy calculation
final_df.groupBy("col1", "col2", "col3", "col4").agg(sum("col5")........sum("coln")).show()

【问题讨论】:

    标签: python apache-spark pyspark apache-spark-sql spark-dataframe


    【解决方案1】:

    Spark SQL 中没有 reduceByKey

    groupBy + 聚合函数的工作方式与 RDD.reduceByKey 几乎相同。 Spark 会自动选择是类似于RDD.groupByKey(即用于collect_list)还是类似于RDD.reduceByKey

    Dataset.groupBy + 聚合函数的性能应该优于或等于 RDD.reduceByKey。 Catalyst 优化器负责如何在后台进行聚合

    【讨论】:

    • 据我记得,它只会在 executor 上添加最终聚合的额外步骤,而不是 Spark SQL groupBy + 聚合中的驱动程序。
    • 感谢您的回复。我们不能在数据帧上应用 reduceByKey 吗?对于大型数据集,许多与 reduceByKey 相同的文章比 group by 更快,因为通过减少最后阶段的行数进行分组。
    • @user3150024 这些文章是关于 RDD 的 - 数据集有一个抽象层,Catalyst 优化器优化查询 :)
    • 还有其他方法可以提高性能吗?试图增加执行器的数量,但它没有反映,只有两个核心用完了 8 个 vcore。我应该将该数据帧转换为 RDD 并应用 reduceByKey。这行得通吗?
    • @user3150024 数据帧的groupByagg 至少应该和reduceByKey 一样快。听起来您还有其他与集群设置相关的问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-11-21
    • 2021-02-10
    • 2020-09-03
    • 2022-08-16
    • 1970-01-01
    • 1970-01-01
    • 2021-11-25
    相关资源
    最近更新 更多