【发布时间】:2017-09-18 15:25:14
【问题描述】:
我想在数据框的所有列上执行窗口函数(具体来说是移动平均)。
我可以这样做
from pyspark.sql import SparkSession, functions as func
df = ...
df.select([func.avg(df[col]).over(windowSpec).alias(col) for col in df.columns])
但恐怕这不是很有效。有没有更好的方法?
【问题讨论】:
-
对我来说似乎不错。
-
感谢您的评论,但这不是窗口操作效率低吗? Spark 会一一对每一列的数据进行分区和排序,还是只执行一次并并行计算所有列的滚动平均值?
-
如果有道理的话,我会讨厌
[partition,sort,calculate_avg(column) for column in columns]这样的东西。 -
您可以查看
DAG以了解幕后情况。 -
我有一个类似的问题,我有几百列,基本上想做一个“填充”类型的操作——对于每一列,对于每一行,如果缺少值,选择它从上一个现有的。您是否设法以任何方式对此进行了优化?
标签: apache-spark apache-spark-sql