【问题标题】:How to use dense_rank and rangeBetween on timestamp value?如何在时间戳值上使用 dense_rank 和 rangeBetween?
【发布时间】: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


    【解决方案1】:

    所以我设法找到适合我的情况的东西,我在这里发布我的答案以防它对某人有所帮助:

    w = Window.partitionBy('userId').orderBy(col("BeginTime"))
    df = df.withColumn('duration_between_series', col('BeginTime').cast('long') - lag(col('EndTime').over(w) )
    .withColumn('rank', dense_rank().over(w))
    .withColumn('sequence_temp', when(col('rank')==1, 1).when(col('duration_between_series')>3600, col('rank')).otherwise(None))
    .withColumn('sequence', last('sequence_temp', True).over(w.rowsBetween(-sys.maxsize, 0))).drop('sequence_temp', 'duration_between_series')
    

    输出:

    +--------------------+--------------------+--------------------+---------+---------+
    |              userId|           BeginTime|             EndTime|     rank| sequence|
    +--------------------+--------------------+--------------------+---------+---------+
    |                   A|2021-02-09 15:56:...|2021-02-09 15:56:...|        1|        1|
    |                   A|2021-02-09 15:57:...|2021-02-09 15:57:...|        2|        1|
    |                   A|2021-02-09 15:58:...|2021-02-09 15:58:...|        3|        1|
    |                   B|2021-02-05 13:16:...|2021-02-05 13:16:...|        1|        1|
    |                   B|2021-02-05 13:16:...|2021-02-05 13:16:...|        2|        1|
    |                   B|2021-02-05 18:27:...|2021-02-05 18:37:...|        3|        3|
    +--------------------+--------------------+--------------------+---------+---------+
    

    列顺序与我预期的不完全一样,但尽管我对每个组都有不同的值,但我很好:)

    【讨论】:

      猜你喜欢
      • 2018-01-13
      • 1970-01-01
      • 2020-05-04
      • 2020-05-24
      • 2015-04-09
      • 1970-01-01
      • 1970-01-01
      • 2021-11-20
      • 1970-01-01
      相关资源
      最近更新 更多