【发布时间】: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