【发布时间】:2017-06-02 09:02:39
【问题描述】:
我有一些数据要按某个列分组,然后根据该组的滚动时间窗口聚合一系列字段。
以下是一些示例数据:
df = spark.createDataFrame([Row(date='2016-01-01', group_by='group1', get_avg=5, get_first=1),
Row(date='2016-01-10', group_by='group1', get_avg=5, get_first=2),
Row(date='2016-02-01', group_by='group2', get_avg=10, get_first=3),
Row(date='2016-02-28', group_by='group2', get_avg=20, get_first=3),
Row(date='2016-02-29', group_by='group2', get_avg=30, get_first=3),
Row(date='2016-04-02', group_by='group2', get_avg=8, get_first=4)])
我想按group_by 进行分组,然后创建从最早日期开始的时间窗口,并一直延伸到该组没有条目的 30 天。在这 30 天结束后,下一个时间窗口将从不属于上一个窗口的下一行的日期开始。
然后我想聚合,例如得到get_avg 的平均值,以及get_first 的第一个结果。
所以这个例子的输出应该是:
group_by first date of window get_avg get_first
group1 2016-01-01 5 1
group2 2016-02-01 20 3
group2 2016-04-02 8 4
编辑:抱歉,我意识到我的问题没有正确指定。我实际上想要一个在 30 天不活动后结束的窗口。我已经相应地修改了示例的 group2 部分。
【问题讨论】:
标签: sql apache-spark pyspark apache-spark-sql window-functions