【问题标题】:Filling pyspark dataframe null values填充pyspark数据框空值
【发布时间】:2020-11-29 16:09:31
【问题描述】:

我有一个 pyspark 数据框,它有 4 列。

+-----+-----+-----+-----+
|col1 |col2 |col3 |col4 |
+-----+-----+-----+-----+
|10   | 5.0 | 5.0 | 5.0 |
|20   | 5.0 | 5.0 | 5.0 |
|null | 5.0 | 5.0 | 5.0 |
|30   | 5.0 | 5.0 | 6.0 |
|40   | 5.0 | 5.0 | 7.0 |
|null | 5.0 | 5.0 | 8.0 |
|50   | 5.0 | 6.0 | 9.0 |
|60   | 5.0 | 7.0 | 10.0|
|null | 5.0 | 8.0 | 11.0|
|70   | 6.0 | 9.0 | 12.0|
|80   | 7.0 | 10.0| 13.0|
|null | 8.0 | 11.0| 14.0|
+-----+-----+-----+-----+

col1 中的某些值缺失,我想根据以下方法设置这些缺失值:

尝试根据具有相同 col2,col3,col4 值的记录的 col1 值的平均值来设置它

如果没有这样的记录,则根据col2,col3值相同的记录中col1的值的平均值来设置

如果还没有这样的记录,则根据具有相同col2值的记录的col1值的平均值设置它

如果以上都找不到,则将其设置为 col1 中所有其他非缺失值的平均值

例如,给定上面的数据帧,只有前两行与第 3 行具有相同的 col2、col3、col4 值。因此,第 3 行的 col1 中的空值应替换为第 1 行中 col1 值的平均值和 2. 对于第 6 行 col1 中的 null 值,它将是第 4 行和第 5 行中 col1 值的平均值,因为只有那些行具有与第 6 行相同的 col2 和 col3 值,而不是相同的 col4 值。和列表继续……

+-----+-----+-----+-----+
|col1 |col2 |col3 |col4 |
+-----+-----+-----+-----+
|10   | 5.0 | 5.0 | 5.0 |
|20   | 5.0 | 5.0 | 5.0 |
|15   | 5.0 | 5.0 | 5.0 |
|30   | 5.0 | 5.0 | 6.0 |
|40   | 5.0 | 5.0 | 7.0 |
|25   | 5.0 | 5.0 | 8.0 |
|50   | 5.0 | 6.0 | 9.0 |
|60   | 5.0 | 7.0 | 10.0|
|35   | 5.0 | 8.0 | 11.0|
|70   | 6.0 | 9.0 | 12.0|
|80   | 7.0 | 10.0| 13.0|
|45   | 8.0 | 11.0| 14.0|
+-----+-----+-----+-----+

最好的方法是什么?

【问题讨论】:

  • For null value in col1 in row 6 你为什么不计算第 1,2,3 行?您的输出似乎不符合您的规则。
  • @Steven 你是对的,我编辑了这个问题。

标签: python dataframe pyspark


【解决方案1】:

我没有找到与您完全相同的值,但是根据您所说的,代码将是这样的:

from pyspark.sql import functions as F

df_2_3_4 = df.groupBy("col2", "col3", "col4").agg(
    F.avg("col1").alias("avg_col1_by_2_3_4")
)
df_2_3 = df.groupBy("col2", "col3").agg(F.avg("col1").alias("avg_col1_by_2_3"))
df_2 = df.groupBy("col2").agg(F.avg("col1").alias("avg_col1_by_2"))
avg_value = df.groupBy().agg(F.avg("col1").alias("avg_col1")).first().avg_col1


df_out = (
    df.join(df_2_3_4, how="left", on=["col2", "col3", "col4"])
    .join(df_2_3, how="left", on=["col2", "col3"])
    .join(df_2, how="left", on=["col2"])
)

df_out.select(
    F.coalesce(
        F.col("col1"),
        F.col("avg_col1_by_2_3_4"),
        F.col("avg_col1_by_2_3"),
        F.col("avg_col1_by_2"),
        F.lit(avg_value),
    ).alias("col1"),
    "col2",
    "col3",
    "col4",
).show()

+----+----+----+----+
|col1|col2|col3|col4|
+----+----+----+----+
|10.0| 5.0| 5.0| 5.0|
|15.0| 5.0| 5.0| 5.0|
|20.0| 5.0| 5.0| 5.0|
|30.0| 5.0| 5.0| 6.0|
|40.0| 5.0| 5.0| 7.0|
|25.0| 5.0| 5.0| 8.0|
|50.0| 5.0| 6.0| 9.0|
|60.0| 5.0| 7.0|10.0|
|35.0| 5.0| 8.0|11.0|
|70.0| 6.0| 9.0|12.0|
|80.0| 7.0|10.0|13.0|
|45.0| 8.0|11.0|14.0|
+----+----+----+----+

【讨论】:

    猜你喜欢
    • 2021-12-07
    • 2016-10-11
    • 2018-01-06
    • 1970-01-01
    • 2023-02-14
    • 2020-07-13
    • 2021-10-21
    • 1970-01-01
    • 2020-08-30
    相关资源
    最近更新 更多