【问题标题】:Filter DataFrame to delete duplicate values in pyspark过滤 DataFrame 以删除 pyspark 中的重复值
【发布时间】:2022-01-13 20:55:15
【问题描述】:

我有以下数据框

date                |   value    | ID
--------------------------------------
2021-12-06 15:00:00      25        1
2021-12-06 15:15:00      35        1
2021-11-30 00:00:00      20        2
2021-11-25 00:00:00      10        2

我想和另一个这样的 DF 一起加入:

idUser | Name | Gender
-------------------
1       John    M
2       Anne    F

我的预期输出是:

ID | Name | Gender | Value
---------------------------
1    John    M        35
2    Anne    F        20

我需要的是:仅获取第一个数据帧的最新值,并仅将此值与我的第二个数据帧连接。虽然,我的 spark 脚本加入了这两个值:

我的代码:

df = df1.select(
   col("date"),
   col("value"),
   col("ID"),
).OrderBy(
   col("ID").asc(),
   col("date").desc(),
).groupBy(
   col("ID"), col("date").cast(StringType()).substr(0,10).alias("date")
).agg (
   max(col("value")).alias("value")
)

final_df = df2.join(
    df,
    (col("idUser") == col("ID")),
    how="left"
)

当我执行这个连接时(格式化列在这篇文章中被抽象出来)我有以下输出:

ID | Name | Gender | Value
---------------------------
1    John    M        35
2    Anne    F        20
2    Anne    F        10

我使用substr 删除小时和分钟以仅按日期过滤。但是当我在不同的日子里有相同的 ID 时,我的输出 df 有 2 个值而不是最近的值。我该如何解决这个问题?

注意:我只使用 pyspark 函数来执行此操作(我现在想使用 spark.sql(...))。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql


    【解决方案1】:

    你可以在pysaprk中使用windowrow_number函数

    from pyspark.sql.window import Window
    from pyspark.sql.functions import row_number
    
    windowSpec = Window.partitionBy("ID").orderBy("date").desc()
    
    df1_latest_val = df1.withColumn("row_number", row_number().over(windowSpec)).filter(
        f.col("row_number") == 1
    )
    

    df1_latest_val 表的输出将如下所示

    date                |   value    | ID | row_number |
    -----------------------------------------------------
    2021-12-06 15:15:00      35        1        1
    2021-11-30 00:00:00      20        2        1
    

    现在您将拥有具有最新 val 的 df,您可以直接将其与另一个表连接。

    【讨论】:

    • Window的作用是什么?
    • 窗口用于对一组数据进行应用操作。就像这里,我们使用ID对数据进行分区,所以ID相同的行将进入一个桶/窗口,然后我们对桶内的数据进行排序,所以在每个桶/窗口中我们都有排序后的一个 ID。然后我们添加行号列并从每个桶/窗口中过滤掉第一行。
    • 我不需要指定我正在使用窗口进行分区的数据框?
    • 首先我们创建了 windowSpec 然后我们在下一行使用了 windowSpec,所以 df 信息被传递给 window 函数。您也可以将两者写在同一行中...
    猜你喜欢
    • 1970-01-01
    • 2020-11-29
    • 2023-03-28
    • 1970-01-01
    • 2016-01-08
    • 2016-05-13
    • 2020-12-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多