使用SparkSession读取CSV文件过滤统计行数时如何处理NullPointerException
问题根因
你遇到的NPE是Spark 2.3.1版本读取CSV时,异常数据行触发的空指针问题,触发场景有两种:
- 部分CSV行首列为空(比如行格式为
,EmployeeX),此时首字段为null,Spark 2.3.1的Row.getString方法读取null值字段时会直接抛出NPE,而非返回null - 存在整行空白的无效CSV行,此时Row长度为0,调用
row.getString(0)会直接抛出异常
你删除filter逻辑后程序正常运行,是因为无需访问Row内部字段,自然不会触发空指针。
解决方案
方案1:修改Row操作逻辑,提前做空值校验
在filter中先校验行长度、字段非空后再取值,避免直接访问null字段:
val validEmployeeIds = Set("12345", "6789") val count = sparkSession .read .option("escape", "\"") .option("quote", "\"") .csv(inputPath) // 先校验行长度合法、首字段非空再判断值是否在合法集合内 .filter(row => row.length >= 1 && !row.isNullAt(0) && validEmployeeIds.contains(row.getString(0))) .distinct() .count() println(count)
方案2:使用Spark DSL语法替代手动Row操作(更推荐)
用Spark内置的列操作API代替手写Row遍历,Spark内部会自动处理空值,不会触发NPE:
val validEmployeeIds = Set("12345", "6789") val count = sparkSession .read .option("escape", "\"") .option("quote", "\"") .csv(inputPath) // 给列指定别名方便后续操作 .toDF("emp_id", "emp_name") // 用内置API判断非空+值是否在合法集合内 .filter(col("emp_id").isNotNull && col("emp_id").isin(validEmployeeIds.toSeq: _*)) .distinct() .count() println(count)
可选优化
- 如果CSV文件包含表头行,可添加
.option("header", "true")避免表头被当成数据行参与计算 - 如果首字段存在大量空字符串,可添加
.option("nullValue", "")统一将空字符串识别为null,配合非空判断过滤即可
内容的提问来源于stack exchange,提问作者Anoop Deshpande
相关产品推荐
相关产品推荐

