【发布时间】:2021-03-08 16:47:22
【问题描述】:
在 Pyspark 中,我尝试使用 dense_rank() 根据 userId 和时间值将行分组到同一组中。
这是我的初始数据框:
+--------------------+--------------------+--------------------+
| userId| BeginTime| EndTime|
+--------------------+--------------------+--------------------+
| A|2021-02-09 15:56:...|2021-02-09 15:56:...|
| A|2021-02-09 15:57:...|2021-02-09 15:57:...|
| A|2021-02-09 15:58:...|2021-02-09 15:58:...|
| B|2021-02-05 13:16:...|2021-02-05 13:16:...|
| B|2021-02-05 13:16:...|2021-02-05 13:16:...|
| B|2021-02-05 18:27:...|2021-02-05 18:37:...|
+--------------------+--------------------+--------------------+
一行代表一个用户执行的一项操作,并给出每个操作的开始日期和结束日期。我想收集连续进行的动作,所以如果两个动作之间的持续时间超过 1 小时,我认为这两个动作不是连续进行的。
这就是我所期望的:
+--------------------+--------------------+--------------------+---------+
| userId| BeginTime| EndTime| sequence|
+--------------------+--------------------+--------------------+---------+
| A|2021-02-09 15:56:...|2021-02-09 15:56:...| 1|
| A|2021-02-09 15:57:...|2021-02-09 15:57:...| 1|
| A|2021-02-09 15:58:...|2021-02-09 15:58:...| 1|
| B|2021-02-05 13:16:...|2021-02-05 13:16:...| 1|
| B|2021-02-05 13:16:...|2021-02-05 13:16:...| 1|
| B|2021-02-05 18:27:...|2021-02-05 18:37:...| 2|
+--------------------+--------------------+--------------------+---------+
我尝试像这样在我的窗口中使用 dense_rank() 和 rangeBetween :
w_rank = (Window
.partitionBy("userId")
.orderBy(col("BeginTime").cast("timestamp").cast("long"))
.rangeBetween(0,3600 )
df = df.withColumn('sequence', dense_rank().over(w_rank))
但我有这个错误:
AnalysisException : Window Frame specifiedwindowframe(RangeFrame, currentrow$(), 3600) must match the required frame specifiedwindowframe(RowFrame, unboundedpreceding$(), currentrow$());
我对 pyspark 很陌生,所以如果有人能在这方面帮助我,我将不胜感激。提前致谢!
【问题讨论】:
标签: pyspark