使用pandas数据创建Spark DataFrame时如何将NaN替换为null
错误原因
PySpark内置的replace()方法不支持将目标值替换为None,因此你之前的写法会触发类型校验错误。
解决方案
你可以选择以下任意一种方法处理NaN值:
方法1:在Spark中用when+isnan函数批量替换
这是最通用的处理方式,针对已生成的Spark DataFrame可以直接修改:
# 先导入依赖函数 from pyspark.sql.functions import col, when, isnan from pyspark.sql.types import DoubleType, IntegerType # 遍历所有数值类型列,批量将NaN替换为null for col_name in spark_df.columns: col_data_type = spark_df.schema[col_name].dataType if isinstance(col_data_type, (DoubleType, IntegerType)): spark_df = spark_df.withColumn( col_name, when(isnan(col(col_name)), None).otherwise(col(col_name)) )
如果只需要处理score单列,也可以简化为单条语句:
spark_df = spark_df.withColumn("score", when(isnan("score"), None).otherwise("score"))
方法2:Pandas读入后先替换NaN再转Spark DataFrame
在创建Spark DataFrame前,先把Pandas里的NaN替换为None,转Spark时会自动映射为null值:
import pandas as pd df = pd.read_csv('data.csv', index_col = False) # pandas中替换NaN为None df = df.where(pd.notnull(df), None) # 后续创建Spark DataFrame的逻辑不变 spark_schema = StructType([ StructField("id",StringType(),True), StructField("segment", StringType(), True), StructField("score",DoubleType(),True), StructField("sales", IntegerType(), True) ]) spark_df = spark.createDataFrame(data=df,schema=spark_schema).cache()
方法3:直接用Spark读取CSV文件
不需要经过Pandas中转,Spark读CSV时默认会将缺失值识别为null,省去转换步骤:
spark_schema = StructType([ StructField("id",StringType(),True), StructField("segment", StringType(), True), StructField("score",DoubleType(),True), StructField("sales", IntegerType(), True) ]) spark_df = spark.read.csv( "data.csv", header=True, # CSV如果有表头保留该参数,没有就设为False schema=spark_schema ).cache()
效果验证
替换完成后执行聚合运算,Spark的avg()函数会自动忽略null值,不会再返回NaN结果。
内容的提问来源于stack exchange,提问作者cs_guy
相关产品推荐
相关产品推荐

