【发布时间】:2018-01-13 20:38:04
【问题描述】:
我的数据如下所示:
userid,eventtime,location_point
4e191908,2017-06-04 03:00:00,18685891
4e191908,2017-06-04 03:04:00,18685891
3136afcb,2017-06-04 03:03:00,18382821
661212dd,2017-06-04 03:06:00,80831484
40e8a7c3,2017-06-04 03:12:00,18825769
如果在同一location_point 的 5 分钟窗口内有 2 个或更多userid,我想添加一个新的布尔列,该列标记为 true。我有一个想法,使用lag 函数在由userid 分区的窗口上查找,范围介于当前时间戳和接下来的 5 分钟之间:
from pyspark.sql import functions as F
from pyspark.sql import Window as W
from pyspark.sql.functions import col
days = lambda i: i * 60*5
windowSpec = W.partitionBy(col("userid")).orderBy(col("eventtime").cast("timestamp").cast("long")).rangeBetween(0, days(5))
lastURN = F.lag(col("location_point"), 1).over(windowSpec)
visitCheck = (last_location_point == output.location_pont)
output.withColumn("visit_check", visitCheck).select("userid","eventtime", "location_pont", "visit_check")
当我使用 RangeBetween 函数时,这段代码给了我一个分析异常:
AnalysisException: u'Window Frame RANGE BETWEEN CURRENT ROW AND 1500 FOLLOWING 必须匹配所需的帧 ROWS BETWEEN 1 PRECEDING 和 1 上一个;
你知道解决这个问题的方法吗?
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql window-functions