【发布时间】:2019-09-30 13:44:43
【问题描述】:
我已经尝试了大约一个月的时间来实现这一点。 仅使用来自其他堆栈溢出问题的一些示例数据。
FinancialYearStart MonthOfFinancialYear SalesTotal
2015 1 10
2015 2 10
2015 5 10
2015 6 50
2016 1 10
2016 3 20
2016 2 30
2017 6 70
2017 7 80
我们如何使用它来计算每个月的年初至今销售额。
FinancialYearStart MonthOfFinancialYear SalesTotal YTDTotal
2015 1 10 10
2015 2 10 20
2015 5 10 30
2015 6 50 50
2016 1 10 60
2016 3 20 80
2016 2 30 110
2017 6 70 70
2017 7 80 150
更具体地说,我们如何使用 group by 或使用窗口函数来计算这些指标。
例如:
Year Month Customer TotalMonthlySales
2015 1 Dog 10
2015 2 Dog 10
2015 3 Cat 20
2015 4 Dog 30
2015 5 Cat 10
2015 7 Cat 20
2015 7 Dog 10
2016 1 Dog 40
2016 2 Dog 20
2016 3 Cat 70
2016 4 Dog 30
2016 5 Cat 10
2016 6 Cat 20
2016 7 Dog 10
愿意:
Year Month Customer TotalMonthlySales YTDSales
2015 1 Dog 10 10
2015 2 Dog 10 20
2015 3 Cat 20 20
2015 4 Dog 30 50
2015 5 Cat 10 30
2015 7 Cat 20 40
2015 7 Dog 10 60
2016 1 Dog 40 40
2016 2 Dog 20 60
2016 3 Cat 70 70
2016 4 Dog 30 90
2016 5 Cat 10 80
2016 6 Cat 20 100
2016 7 Dog 10 100
现在我使用滚动窗口函数使用以下方法计算过去 4 周 13 周。
w = (Window().partitionBy(col("number"),col("id")).orderBy(F.col("timestampGMT").cast('long')).rangeBetween(-days(27), 0))
df1 = df.withColumn('lw_4weeksCY_SALES_DOLLARS', F.sum("CY_SALES_DOLLARS").over(w)).withColumn('lw_4weeksCY_SALES_UNITS', F.sum("CY_SALES_UNITS")
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql