【发布时间】:2016-11-19 18:22:05
【问题描述】:
我有一个这样的数据框:
id x y
1 a 1 P
2 a 2 S
3 b 3 P
4 b 4 S
我想保留 y 的“前导”值为“S”的行让我们说,这样我的结果数据框将是:
id x y
1 a 1 P
2 b 3 P
我可以使用 PySpark 执行以下操作:
getLeadPoint = udf(lambda x: 'S' if (y == 'S') else 'NOTS', StringType())
windowSpec = Window.partitionBy(df['id'])
df = df.withColumn('lead_point', getLeadPoint(lead(df.y).over(windowSpec)))
dfNew = df.filter(df.lead_point == 'S')
但是,在这里,我正在改变一个不必要的列,然后进行过滤。
我想要做的是这样的事情,我使用 Lead() 进行过滤,但无法使其工作:
dfNew = df.filter(lead(df.y).over(windowSpec) == 'S')
关于如何使用窗口直接过滤来实现结果的任何想法?
R 等价物是:
library(dplyr)
df %>% group_by(id) %>% filter(lead(y) == 'S')
【问题讨论】:
-
对不起....顺便说一句 - 订购部分是一个占位符。我需要订购的数据中有一个“时间戳”列。
-
这里其实有个bug。最简单的解决方法是使用
withColumn添加一列并将其用于过滤。 -
感谢您的验证。那么,我目前的解决方案是我能做的最好的吗?
-
这里不需要
udf,只需过滤lead()的输出即可。
标签: apache-spark pyspark apache-spark-sql pyspark-sql window-functions