【问题标题】:Fill Pyspark dataframe column null values with average value from same column用同一列的平均值填充 Pyspark 数据框列空值
【发布时间】:2016-10-11 12:16:19
【问题描述】:

有了这样的数据框,

rdd_2 = sc.parallelize([(0,10,223,"201601"), (0,10,83,"2016032"),(1,20,None,"201602"),(1,20,3003,"201601"), (1,20,None,"201603"), (2,40, 2321,"201601"), (2,30, 10,"201602"),(2,61, None,"201601")])

df_data = sqlContext.createDataFrame(rdd_2, ["id", "type", "cost", "date"])
df_data.show()

+---+----+----+-------+
| id|type|cost|   date|
+---+----+----+-------+
|  0|  10| 223| 201601|
|  0|  10|  83|2016032|
|  1|  20|null| 201602|
|  1|  20|3003| 201601|
|  1|  20|null| 201603|
|  2|  40|2321| 201601|
|  2|  30|  10| 201602|
|  2|  61|null| 201601|
+---+----+----+-------+

我需要用现有值的平均值填充空值,预期结果为

+---+----+----+-------+
| id|type|cost|   date|
+---+----+----+-------+
|  0|  10| 223| 201601|
|  0|  10|  83|2016032|
|  1|  20|1128| 201602|
|  1|  20|3003| 201601|
|  1|  20|1128| 201603|
|  2|  40|2321| 201601|
|  2|  30|  10| 201602|
|  2|  61|1128| 201601|
+---+----+----+-------+

其中1128 是现有值的平均值。我需要为几列这样做。

我目前的做法是使用na.fill:

fill_values = {column: df_data.agg({column:"mean"}).flatMap(list).collect()[0] for column in df_data.columns if column not in ['date','id']}
df_data = df_data.na.fill(fill_values)

+---+----+----+-------+
| id|type|cost|   date|
+---+----+----+-------+
|  0|  10| 223| 201601|
|  0|  10|  83|2016032|
|  1|  20|1128| 201602|
|  1|  20|3003| 201601|
|  1|  20|1128| 201603|
|  2|  40|2321| 201601|
|  2|  30|  10| 201602|
|  2|  61|1128| 201601|
+---+----+----+-------+

但这很麻烦。有什么想法吗?

【问题讨论】:

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


    【解决方案1】:

    好吧,你必须以一种或另一种方式:

    • 计算统计数据
    • 填空

    它几乎限制了你在这里可以真正改进的地方,仍然:

    • flatMap(list).collect()[0]替换为first()[0]或结构解包
    • 通过单个操作计算所有统计数据
    • 使用内置的Row方法提取字典

    最终的结果可能是这样的:

    def fill_with_mean(df, exclude=set()): 
        stats = df.agg(*(
            avg(c).alias(c) for c in df.columns if c not in exclude
        ))
        return df.na.fill(stats.first().asDict())
    
    fill_with_mean(df_data, ["id", "date"])
    

    在 Spark 2.2 或更高版本中,您还可以使用 Imputer。见Replace missing values with mean - Spark Dataframe

    【讨论】:

    • @Kevad 将 pyspark.sql.functions 导入为 fn,然后使用 fn.avg
    • 有人改进了这个功能吗?占用我太多的计算时间!? :)
    • @NicoCoallier 性能方面,您不会得到任何显着改善。这几乎是最佳解决方案。 API 方面,您可以使用 Imputer
    • 排除还是包含更快?
    • @NicoCoallier 这影响可以忽略不计,但如果您在大量列(如数千列)上执行此操作,仅优化器开销就很重要。如果是这种情况,RDD 可能会更快。当然,很大程度上取决于您所说的慢...
    猜你喜欢
    • 2021-12-07
    • 1970-01-01
    • 2020-11-29
    • 2017-04-14
    • 2019-11-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多