Lik*_*iky 4 python dataframe pyspark
我有一个大约的数据帧。400 万行和 35 列作为输入。
我对此数据框所做的只是以下步骤:
因此,我们最终得到与我们开始时相同的数据帧(理论上)。
但是,我注意到,如果给定列的列表变得太大(超过 6 列),则输出数据帧将变得无法操作。即使是简单的显示也需要 10 分钟。
这是我的代码示例(df 是我的输入数据帧):
for c in list_columns:
df = df.join(df.groupby(list_group_features).agg(sum(c).alias('sum_' + c)), list_group_features)
df = df.drop('sum_' + c)
Run Code Online (Sandbox Code Playgroud)
发生这种情况是由于 Spark 的内部工作原理及其惰性求值造成的。
groupby当您调用、join、时, Spark 会执行什么操作agg,它将这些调用附加到对象的计划中df。因此,即使它没有对数据执行任何操作,您也会创建一个大型执行计划,该计划内部存储在 Spark DataFrame 对象中。
仅当您调用操作(show、count、write等)时,Spark 才会优化计划并执行它。如果计划太大,优化步骤可能需要一段时间才能执行。还要记住,计划优化发生在驱动程序上,而不是执行程序上。因此,如果您的驱动程序繁忙或超载,它也会延迟 Spark 计划优化步骤。
记住,无论对于优化还是执行来说,连接在 Spark 中都是昂贵的操作。如果可以的话,在操作单个 DataFrame 时应始终避免连接,而应使用窗口功能。仅当您连接来自不同源(不同表)的不同数据帧时才应使用连接。
优化代码的一种方法是:
import pyspark
import pyspark.sql.functions as f
w = pyspark.sql.Window.partitionBy(list_group_features)
agg_sum_exprs = [f.sum(f.col(c)).alias("sum_" + c).over(w) for c in list_columns]
res_df = df.select(df.columns + agg_sum_exprs)
Run Code Online (Sandbox Code Playgroud)
list_group_features对于大型列表来说,这应该是可扩展且快速的list_columns。