【问题标题】:Filter rows where the lead/lag are specific values (window with filter)过滤超前/滞后是特定值的行(带过滤器的窗口)
【发布时间】: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


【解决方案1】:

效率不高,但你可以用索引压缩,然后创建一个新的RDD,在索引上加1,然后加入索引,然后变成一个简单的过滤操作。

【讨论】:

  • 谢谢!我认为这是一个比我改变一个新的潜在客户值列更糟糕的解决方案,就像我目前正在做的那样。此外,如果没有代码,很难说出您的解决方案实际上是什么样子或做了什么。
【解决方案2】:

假设您的数据如下所示:

df = sc.parallelize([
    ("a", 1,  1, "P"), ("a", 2,  2, "S"),
    ("b", 4,  2, "S"), ("b", 3,  1, "P"), ("b", 2,  3, "P"), ("b", 3,  3, "S")
]).toDF(["id", "x", "timestamp", "y"])

window spec 等价于

from pyspark.sql.functions import lead, col
from pyspark.sql import Window

w = Window.partitionBy("id").orderBy("timestamp")

您可以简单地添加列并将其用于过滤:

(df
    .withColumn("lead_y", lead("y").over(w))
    .where(col("lead_y") == "S").drop("lead_y"))

它并不漂亮,但会比 UDF 调用更有效。

【讨论】:

    猜你喜欢
    • 2022-12-13
    • 1970-01-01
    • 2018-06-09
    • 1970-01-01
    • 1970-01-01
    • 2019-01-02
    • 1970-01-01
    • 2013-07-20
    • 2011-11-01
    相关资源
    最近更新 更多