【问题标题】:How to calculate date difference in pyspark using window function?如何使用窗口函数计算 pyspark 中的日期差异?
【发布时间】:2020-04-03 05:25:20
【问题描述】:

尝试计算自用户首次开始使用应用程序以来经过的天数以及 df 行所代表的事件。下面的代码 (via) 创建了一个列,将行与前一行进行比较,但我需要将它与分区的第一行进行比较。

window = Window.partitionBy('userId').orderBy('dateTime')

df = df.withColumn("daysPassed", datediff(df.dateTime, 
                                lag(df.dateTime, 1).over(window)))

尝试用“int(Window.unboundedPreceding)”代替 1,结果报错。

我希望 daysPassed 列执行的操作示例:

 Row(userId='59', page='NextSong', datetime='2018-10-01', daysPassed=0),
 Row(userId='59', page='NextSong', datetime='2018-10-03', daysPassed=2),
 Row(userId='59', page='NextSong', datetime='2018-10-04', daysPassed=3)

【问题讨论】:

  • 将lag()更改为min()
  • 获取:TypeError:列不可迭代
  • 可能是与python的min()发生冲突。最好使用模块参考

标签: python apache-spark pyspark


【解决方案1】:

所以,如果我做对了,基本上你会想要计算行中的日期与用户的最小日期(开始日期)的差异,而不是 lag()。

from pyspark.sql import functions as func
window = Window.partitionBy('userId')

df_b = df_a.withColumn("daysPassed", func.datediff(df.dateTime, func.min(df.dateTime).over(window)))

这会计算用户第一次启动应用程序的天数。

【讨论】:

  • 获取:TypeError:列不可迭代
  • 使用from pyspark.sql import functions as func,然后使用... func.min(df.dateTime).over(window)))
  • 谢谢。不幸的是,现在我遇到了另一个错误。我会做一些挖掘。错误:AnalysisException:“分组表达式序列为空,并且'artist'不是聚合函数。在窗口函数中包装'(min(dateTime) AS _w0)'或包装'artist ' 在 first() (或 first_value)中,如果您不在乎获得哪个值。
猜你喜欢
  • 2017-10-16
  • 2019-04-07
  • 1970-01-01
  • 2016-08-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-07
  • 2014-09-15
相关资源
最近更新 更多