【问题标题】:Apply window function over multiple columns在多列上应用窗口函数
【发布时间】: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


【解决方案1】:

另一种可能更好的方法是创建一个新的 df,在其中您按 Window 函数中的列分组并在其余列上应用平均值,然后进行左连接。对于 df 溢出到磁盘(或无法持久保存在内存中)的大型数据帧,这肯定会更优化。

【讨论】:

  • 在上述情况下,窗口函数会将 df 溢出到磁盘吗?我看不出 groupby 列是否与 df 已被分区的相同
猜你喜欢
  • 1970-01-01
  • 2010-12-26
  • 1970-01-01
  • 1970-01-01
  • 2016-04-06
  • 2017-08-02
  • 2017-11-09
  • 2018-01-16
  • 1970-01-01
相关资源
最近更新 更多