【问题标题】:How to calculate cumulative sum over date range excluding weekends in PySpark 2.0?如何在 PySpark 2.0 中计算不包括周末的日期范围内的累积总和?
【发布时间】:2020-12-04 23:04:49
【问题描述】:

这是对我之前在这里提出的问题How to calculate difference between dates excluding weekends in PySpark 2.2.0 的扩展。我的 spark 数据框如下所示,可以使用随附的代码生成:

df = spark.createDataFrame([(1, "John Doe", "2020-11-30",1),(2, "John Doe", "2020-11-27",2),(4, "John Doe", "2020-12-01",0),(5, "John Doe", "2020-10-02",1),\
                          (6, "John Doe", "2020-12-03",1),(7, "John Doe", "2020-12-04",1)],
                            ("id", "name", "date","count"))

+---+--------+----------+-----+
| id|    name|      date|count|
+---+--------+----------+-----+
|  5|John Doe|2020-10-02|    1|
|  2|John Doe|2020-11-27|    2|
|  1|John Doe|2020-11-30|    1|
|  4|John Doe|2020-12-01|    0|
|  6|John Doe|2020-12-03|    1|
|  7|John Doe|2020-12-04|    1|
+---+--------+----------+-----+

我正在尝试计算 2、3、4、5 和 30 天期间的累积总和。下面是2天的示例代码和结果表。

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf
days = lambda i: i * 86400
windowval_2 = Window.partitionBy("name").orderBy(F.col("date").cast("timestamp").cast("long")).rangeBetween(days(-1), days(0))
windowval_3 = Window.partitionBy("name").orderBy(F.col("date").cast("timestamp").cast("long")).rangeBetween(days(-2), days(0))
windowval_4 = Window.partitionBy("name").orderBy(F.col("date").cast("timestamp").cast("long")).rangeBetween(days(-3), days(0))
df = df.withColumn("cum_sum_2d_temp",F.sum("count").over(windowval_2))


+---+--------+----------+-----+---------------+
| id|    name|      date|count|cum_sum_2d_temp|
+---+--------+----------+-----+---------------+
|  5|John Doe|2020-10-02|    1|              1|
|  2|John Doe|2020-11-27|    2|              2|
|  1|John Doe|2020-11-30|    1|              1|
|  4|John Doe|2020-12-01|    0|              1|
|  6|John Doe|2020-12-03|    1|              1|
|  7|John Doe|2020-12-04|    1|              2|
+---+--------+----------+-----+---------------+

我想做的是在计算日期范围时,计算不包括周末,即在我的表中 2020-11-27 是星期五,2020-11-30 是星期一。如果我们排除周六和周日,它们之间的差异为 1。我希望 2020-11-27 和 2020-11-30 值的累积总和在“cum_sum_2d_temp”列中的 2020-11-30 之前应该是 3。我希望将解决方案结合到我之前的问题到日期范围。

【问题讨论】:

  • 您是否尝试过使用 filter 和自定义谓词函数来检查您的特定需求?
  • @OakenDuck 我使用df.withColumn("sum_2d",F.when(workdaysUDF(F.col("date"),F.lag(F.col("date"),1).over(windowVal_gen))==1,\ F.col("sum_2d_temp")).otherwise(F.col("count"))) 来检查两个日期之间的 datediff,但这种方法的缺点是,由于我没有每个 id 的每日数据,我不知道要连续多少个日期检查每个cum sum。正如我在 OP 中所说,我希望对 2、3、4 和 5 天求和,因此两个连续行之间的 datediff 可以是任何值,例如 3、4 等。
  • @OakenDuck 抱歉,我错过了为工作日添加 UDF workdaysUDF = F.udf(lambda date1, date2: int(np.busday_count(date2, date1)) if (date1 is not None and date2 is not None) else None, IntegerType())

标签: python apache-spark pyspark


【解决方案1】:

计算相对于最早日期的 date_dif:

import numpy as np
import pyspark.sql.functions as F
from pyspark.sql.window import Window
from pyspark.sql.types import IntegerType

df = spark.createDataFrame([(1, "John Doe", "2020-11-30",1),(2, "John Doe", "2020-11-27",2),(4, "John Doe", "2020-12-01",0),(5, "John Doe", "2020-10-02",1),\
                          (6, "John Doe", "2020-12-03",1),(7, "John Doe", "2020-12-04",1)],
                            ("id", "name", "date","count"))

workdaysUDF = F.udf(lambda date1, date2: int(np.busday_count(date2, date1)) if (date1 is not None and date2 is not None) else None, IntegerType())
df = df.withColumn("date_dif", workdaysUDF(F.col('date'), F.first(F.col('date')).over(Window.partitionBy('name').orderBy('date'))))

windowval = lambda days: Window.partitionBy('name').orderBy('date_dif').rangeBetween(-days, 0) 
df = df.withColumn("cum_sum",F.sum("count").over(windowval(2)))
df.show()

+---+--------+----------+-----+--------+-------+
| id|    name|      date|count|date_dif|cum_sum|
+---+--------+----------+-----+--------+-------+
|  5|John Doe|2020-10-02|    1|       0|      1|
|  2|John Doe|2020-11-27|    2|      40|      2|
|  1|John Doe|2020-11-30|    1|      41|      3|
|  4|John Doe|2020-12-01|    0|      42|      3|
|  6|John Doe|2020-12-03|    1|      44|      1|
|  7|John Doe|2020-12-04|    1|      45|      2|
+---+--------+----------+-----+--------+-------+

【讨论】:

  • 谢谢,这对我来说非常有效。但是,您能否详细说明一下 windowval lambda 函数的作用。
  • @vagautam 它只是使窗口中的天数变为动态
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-12-12
  • 1970-01-01
  • 2018-03-05
  • 1970-01-01
  • 1970-01-01
  • 2023-03-06
相关资源
最近更新 更多