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