PySpark用均值填充缺失值出现Py4JJavaError异常如何解决
问题原因
- 核心问题1:CSV读取未开启类型推断,所有列默认被解析为字符串类型。
avg()函数对字符串列计算均值时会直接返回null,而na.fill()方法不允许填充值为null,就会触发本次的NullPointerException。你代码中虽然定义了infer_schema = "true"变量,但没有传入读取配置的option中,配置未生效。 - 核心问题2:参数逻辑错误。你把需要填充均值的
MinTemp、MaxTemp、Evaporation、Sunshine列传入了exclude参数,会导致这些列直接被排除在均值计算范围外,就算类型正确也无法拿到对应填充值。 - 潜在问题:如果某列所有值均为空,
avg()返回结果也为null,同样会触发空指针错误,需要额外过滤这种情况。
修复方案
第一步:修复CSV读取逻辑
开启类型推断,确保数值列被解析为数字类型:
df_1= spark.read.format("csv") \ .option("header","true") \ .option("inferSchema","true") \ .load('/content/weatherAUS.csv')
第二步:修复均值填充函数
通用版本:填充所有数值列的空值
仅对数值类型列计算均值,自动过滤均值为空的列,避免空指针:
from pyspark.sql.functions import avg def fill_with_mean(df, exclude=set()): # 仅筛选数值类型、且不在排除列表的列做均值计算 numeric_cols = [c for c, t in df.dtypes if t in ('int', 'double', 'float') and c not in exclude] # 计算各列均值 avg_result = df.agg(*(avg(c).alias(c) for c in numeric_cols)).first().asDict() # 过滤掉均值为null的列,避免空指针 fill_map = {k: v for k, v in avg_result.items() if v is not None} return df.na.fill(fill_map) # 示例:排除非数值的分类列、日期列,填充其余所有数值列 res = fill_with_mean(df_1, exclude=["Date", "Location", "WindGustDir", "RainToday", "RainTomorrow"]) res.show()
定制版本:仅填充指定列的均值
如果你只需要填充MinTemp、MaxTemp、Evaporation、Sunshine四列,可直接用这个简化版本:
from pyspark.sql.functions import avg def fill_target_cols_mean(df, target_cols): avg_result = df.agg(*(avg(c).alias(c) for c in target_cols)).first().asDict() fill_map = {k: v for k, v in avg_result.items() if v is not None} return df.na.fill(fill_map) res = fill_target_cols_mean(df_1, target_cols=["MinTemp", "MaxTemp", "Evaporation", "Sunshine"]) res.show()
内容的提问来源于stack exchange,提问作者John
相关产品推荐
相关产品推荐

