【发布时间】:2017-09-29 20:53:22
【问题描述】:
我有一个“user_name”、“mac”、“dayte”(天)的数据集。我想按 ['user_name'] 分组。然后为该 GROUP BY 创建使用“dayte”滚动 30 天的窗口。在滚动的 30 天期间,我想计算“mac”的不同数量。并将其添加到我的数据框中。数据样本。
user_name mac dayte
0 001j 7C:D1 2017-09-15
1 0039711 40:33 2017-07-25
2 0459 F0:79 2017-08-01
3 0459 F0:79 2017-08-06
4 0459 F0:79 2017-08-31
5 0459 78:D7 2017-09-08
6 0459 E0:C7 2017-09-16
7 133833 18:5E 2017-07-27
8 133833 F4:0F 2017-07-31
9 133833 A4:E4 2017-08-07
我已尝试使用 PANDAs 数据框解决此问题。
df['ct_macs'] = df.groupby(['user_name']).rolling('30d', on='dayte').mac.apply(lambda x:len(x.unique()))
但收到错误
Exception: cannot handle a non-unique multi-index!
我在 PySpark 中尝试过,但也收到了错误。
from pyspark.sql import functions as F
#function to calculate number of seconds from number of days
days = lambda i: i * 86400
#convert string timestamp to timestamp type
df= df.withColumn('dayte', df.dayte.cast('timestamp'))
#create window by casting timestamp to long (number of seconds)
w = Window.partitionBy("user_name").orderBy("dayte").rangeBetween(-days(30), 0)
df= df.select("user_name","mac","dayte",F.size(F.denseRank().over(w).alias("ct_mac")))
但收到错误
Py4JJavaError: An error occurred while calling o464.select.
: org.apache.spark.sql.AnalysisException: Window function dense_rank does not take a frame specification.;
我也试过
df= df.select("user_name","dayte",F.countDistinct(col("mac")).over(w).alias("ct_mac"))
但显然,Spark 不支持它(在 Window 中计数不同)。 我对纯粹的 SQL 方法持开放态度。在 MySQL 或 SQL Server 中,但更喜欢 Python 或 Spark。
【问题讨论】:
标签: python sql pandas pyspark pyspark-sql