【问题标题】:Replace Null values with median in pyspark用pyspark中的中位数替换空值
【发布时间】:2019-01-11 19:26:04
【问题描述】:

如何在数据集 df 下方的 Age 和 Height 列中用中位数替换空值。

df = spark.createDataFrame([(1, 'John', 1.79, 28,'M', 'Doctor'),
                        (2, 'Steve', 1.78, 45,'M', None),
                        (3, 'Emma', 1.75, None, None, None),
                        (4, 'Ashley',1.6, 33,'F', 'Analyst'),
                        (5, 'Olivia', 1.8, 54,'F', 'Teacher'),
                        (6, 'Hannah', 1.82, None, 'F', None),
                        (7, 'William',None, 42,'M', 'Engineer'),
                        (None,None,None,None,None,None),
                        (8,'Ethan',1.55,38,'M','Doctor'),
                        (9,'Hannah',1.65,None,'F','Doctor'),
                       (10,'Xavier',1.64,43,None,'Doctor')]
                       , ['Id', 'Name', 'Height', 'Age', 'Gender', 'Profession'])

在Replace missing values with mean - Spark Dataframe的帖子中,我使用了给出的函数 from pyspark.ml.feature 导入 Imputer

imputer = Imputer(
inputCols=df.columns, 
outputCols=["{}_imputed".format(c) for c in df.columns])

imputer.fit(df).transform(df)

它给我一个错误。

IllegalArgumentException: '要求失败:列 ID 的类型必须等于以下类型之一:[DoubleType, FloatType] 但实际上是 LongType 类型。'

所以请帮忙。 谢谢

【问题讨论】:

标签: replace null pyspark median


【解决方案1】:

这可能是初始转换错误(我有一些字符串需要浮动)。要将所有 cols 转换为浮点数,请执行以下操作:

from pyspark.sql.functions import col
df = df.select(*(col(c).cast("float").alias(c) for c in df.columns))

那么你应该可以估算。注意:我将我的策略设置为中值而不是平均值。

from pyspark.ml.feature import Imputer

imputer = Imputer(
    inputCols=df.columns, 
    outputCols=["{}_imputed".format(c) for c in df.columns]
    ).setStrategy("median")

# Add imputation cols to df
df = imputer.fit(df).transform(df)

【讨论】:

  • 您还可以使用outputCols=[c for c in df.columns] 的估算值覆盖当前列名。另请注意,由于工作节点之间的分布,有时median 很难让 PySpark 计算。似乎mean 通常更稳定。 This question 提供了一个潜在的解决方案。
【解决方案2】:

我对更优雅的解决方案感兴趣,但我分别从数字中估算了分类。为了估算分类值,我得到了最常见的值,并使用 when 和 otherwise 函数填空:

import pyspark.sql.functions as F
for col_name in ['Name', 'Gender', 'Profession']:
    common = df.dropna().groupBy(col_name).agg(F.count("*")).orderBy('count(1)', ascending=False).first()[col_name]
    df = df.withColumn(col_name, F.when(F.isnull(col_name), common).otherwise(df[col_name]))

为了在运行 imputer 行之前估算数字,我只是将 Age 和 Id 列转换为双精度,从而规避了数字字段的问题,并将 imputer 限制为数字列。

from pyspark.ml.feature import Imputer
df = df.withColumn("Age", df['Age'].cast('double')).withColumn('Id', df['Id'].cast('double'))
imputer = Imputer(
inputCols=['Id', 'Height', 'Age'],
outputCols=['Id', 'Height', 'Age'])
imputer.fit(df).transform(df)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-10-23
    • 1970-01-01
    • 1970-01-01
    • 2020-04-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-09
    相关资源
    最近更新 更多