You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用SparkSession读取CSV文件过滤统计行数时如何处理NullPointerException

问题根因

你遇到的NPE是Spark 2.3.1版本读取CSV时,异常数据行触发的空指针问题,触发场景有两种:

  1. 部分CSV行首列为空(比如行格式为,EmployeeX),此时首字段为null,Spark 2.3.1的Row.getString方法读取null值字段时会直接抛出NPE,而非返回null
  2. 存在整行空白的无效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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 19:36:02