【问题标题】:PySpark window function to get last row with date column value equal to datePySpark窗口函数获取日期列值等于日期的最后一行
【发布时间】:2020-05-15 21:45:34
【问题描述】:

我试图让一个窗口函数返回并在特定日期之前获取上一行,我不太确定出了什么问题,但它给了我上一行而不是指定的日期行。为了计算这一点,我正在获取当前行日期并找到与该周相关的当前星期一

    def previous_day(date, dayOfWeek):
        return date_sub(next_day(date, "monday"), 7)
    spark_df = spark_df.withColumn("last_monday", previous_day(spark_df['calendarday'], "monday"))

然后我正在计算当天与最近的前一个星期一之间的差异(以天为单位)

    d = F.datediff(spark_df['calendarday'], spark_df['last_monday'])
    spark_df = spark_df.withColumn("daysSinceMonday",d)

从我的 daysSinceMonday 中可以看出,每行的值都是正确的。接下来我想创建一个窗口并选择它的第一行,但将它们按我设置的 d 值进行范围,但由于某种原因它不起作用。

    days = lambda i: i * 86400 
    w = (Window.partitionBy(column_list).orderBy(col('calendarday').cast("timestamp").cast("long")).rangeBetween(-days(d), 0))
    spark_df = spark_df.withColumn('PreviousYearUnique', first("indexCP").over(w))

    Starting Data Frame
    ## +---+-----------+-----------+--------+       
    ## | id|calendarday|last_monday| indexCP|
    ## +---+-----------+-----------+--------+
    ## |  1|2015-01-05 | 2015-01-05|  0.0076|
    ## |  1|2015-01-06 | 2015-01-05|  0.0026|
    ## |  1|2015-01-07 | 2015-01-05|  0.0016|
    ## |  1|2015-01-08 | 2015-01-05|  0.0006|
    ## |  2|2015-01-09 | 2015-01-05|  0.0012|
    ## |  2|2015-01-10 | 2015-01-05|  0.0014|
    ## |  1|2015-01-12 | 2015-01-12|  0.0026|
    ## |  1|2015-01-13 | 2015-01-12|  0.0086|
    ## |  1|2015-01-14 | 2015-01-12|  0.0046|
    ## |  1|2015-01-15 | 2015-01-12|  0.0021|
    ## |  2|2015-01-16 | 2015-01-12|  0.0042|
    ## |  2|2015-01-17 | 2015-01-12|  0.0099|
    ## +---+-----------+-----------+--------+

    New Data Frame Adding Previous last_mondays row indexCP as PreviousYearUnique
    ## +---+-----------+-----------+--------+--------------------+       
    ## | id|calendarday|last_monday| indexCP| PreviousYearUnique |
    ## +---+-----------+-----------+--------+--------------------+
    ## |  1|2015-01-05 | 2015-01-05|  0.0076|              0.0076|
    ## |  1|2015-01-06 | 2015-01-05|  0.0026|              0.0076|
    ## |  1|2015-01-07 | 2015-01-05|  0.0016|              0.0076|
    ## |  1|2015-01-08 | 2015-01-05|  0.0006|              0.0076|
    ## |  2|2015-01-09 | 2015-01-05|  0.0012|              0.0076|
    ## |  2|2015-01-10 | 2015-01-05|  0.0014|              0.0076|
    ## |  1|2015-01-12 | 2015-01-12|  0.0026|              0.0026|
    ## |  1|2015-01-13 | 2015-01-12|  0.0086|              0.0026|
    ## |  1|2015-01-14 | 2015-01-12|  0.0046|              0.0026|
    ## |  1|2015-01-15 | 2015-01-12|  0.0021|              0.0026|
    ## |  2|2015-01-16 | 2015-01-12|  0.0042|              0.0026|
    ## |  2|2015-01-17 | 2015-01-12|  0.0099|              0.0026|
    ## +---+-----------+-----------+--------+--------------------+

有什么想法吗?

【问题讨论】:

  • 如果您可以以表格格式提供示例数据和所需的输出,这将有助于人们回答。(欢迎使用 SO)
  • 好点添加了它们。感谢这是一个整洁的地方!

标签: python apache-spark pyspark aws-glue


【解决方案1】:

您可以在 unboundedPreceding 窗口上通过 calendarday partitionBy last_monday,然后使用 first

from pyspark.sql import functions as F
from pyspark.sql.window import Window

w=Window().partitionBy("last_monday")\
          .orderBy(F.to_date("calendarday","yyyy-MM-dd"))\
          .rowsBetween(Window.unboundedPreceding,Window.currentRow)

df.withColumn("PreviousYearUnique", F.first("indexCP").over(w)).show()


#+---+-----------+-----------+-------+------------------+
#| id|calendarday|last_monday|indexCP|PreviousYearUnique|
#+---+-----------+-----------+-------+------------------+
#|  1| 2015-01-05| 2015-01-05| 0.0076|            0.0076|
#|  1| 2015-01-06| 2015-01-05| 0.0026|            0.0076|
#|  1| 2015-01-07| 2015-01-05| 0.0016|            0.0076|
#|  1| 2015-01-08| 2015-01-05| 6.0E-4|            0.0076|
#|  2| 2015-01-09| 2015-01-05| 0.0012|            0.0076|
#|  2| 2015-01-10| 2015-01-05| 0.0014|            0.0076|
#|  1| 2015-01-12| 2015-01-12| 0.0026|            0.0026|
#|  1| 2015-01-13| 2015-01-12| 0.0086|            0.0026|
#|  1| 2015-01-14| 2015-01-12| 0.0046|            0.0026|
#|  1| 2015-01-15| 2015-01-12| 0.0021|            0.0026|
#|  2| 2015-01-16| 2015-01-12| 0.0042|            0.0026|
#|  2| 2015-01-17| 2015-01-12| 0.0099|            0.0026|
#+---+-----------+-----------+-------+------------------+

【讨论】:

  • 数据共有 40 列,为了便于查看,我只是将其删减。它目前必须是partitionBy(column_list),它被定义为column_list = ["accountname","secname"]
  • 当然,你必须将上周一添加到 column_list,然后使用显示的窗口首先计算
  • 好的,我会试一试,但我认为对于 id: 1 calendarday: 2015-01-15 last_monday: 2015-01-12 它会给我#| 1| 2015-01-05| 2015-01-05| 0.0076| 0.0076| 但我需要的行是#| 1| 2015-01-12| 2015-01-12| 0.0026| 0.0026|
  • 看起来确实有效。我不明白为什么它有效,但它确实有效!谢谢!
猜你喜欢
  • 1970-01-01
  • 2021-01-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-12-14
  • 2019-02-15
  • 2017-11-22
  • 2019-02-03
相关资源
最近更新 更多