Scala调用collect()函数时出现NullPointerException问题求助
咱们一步步拆解这个问题,先搞清楚为什么take(6)能正常运行,换成collect()就抛出空指针异常:
第一个核心原因:空值dataFormat直接触发NPE
看你的source.csv数据,里面有几行的Format列是空的(比如EMP_NAME、HEADER_ROW_COUNT)。当这些行进入_getHiveTypeMapping函数时,dataFormat参数是null,这时候你直接调用dataFormat.toUpperCase.contains("DATETIME"),必然会触发空指针异常!
那为什么take(6)没报错?因为Spark在本地模式下,小数据量的take操作可能直接在Driver端执行,刚好没走到那些空值的逻辑分支;而collect()会遍历所有行,必然会碰到这些空值行,问题就暴露出来了。
第二个关键问题:SparkContext在Executor中无法序列化
你的函数里每次调用都用sc.textFile加载映射文件:
val mappings = sc.textFile(sasToHiveMappingLocation)
sc是Driver端创建的SparkContext对象,当函数被序列化到Executor节点执行时,sc根本没法被正确序列化到Executor环境中,到了Executor那边sc就是null,调用sc.textFile自然就抛NPE了。
还是那个逻辑:take(6)数据量小,可能全程在Driver端执行,不需要把sc传到Executor,所以没暴露这个问题;collect()要分发任务到Executor,这个隐藏的问题就立刻显现了。
额外的坑:映射文件加载方式既低效又不安全
每次调用函数都重新读一遍映射文件,不仅浪费IO资源、拖慢性能,还会因为重复加载增加出错的概率,完全没必要。
修复方案,一步一步来
1. 先处理空值dataFormat
在使用dataFormat之前必须判空,修改函数里的逻辑:
if(dataFormat != null && dataFormat.toUpperCase.contains("DATETIME")){ definedType="datetime" } else if(dataFormat != null) { // 先确认非空再处理数字判断 try { // 后面会修正数字判断逻辑,先暂时保留结构 if(isNumeric(dataFormat)) { definedType="Double" } } catch { case _: Throwable => definedType="Unknown" } } else { definedType="Unknown" }
2. 提前加载映射文件并广播
在Driver端提前把映射文件加载成Map结构,然后广播到所有Executor,这样每个任务都能复用这个映射,不用重复读文件:
// 在Driver端执行,提前加载映射并转成Map val sasToHiveMappingLocation = "s3a://abc/SASToHiveMappings.csv" val mappingsMap = sc.textFile(sasToHiveMappingLocation) .map(line => line.split(",")) .filter(arr => arr.length >=3) // 过滤非法行 .map(arr => ((arr(0).toUpperCase, arr(1).toUpperCase), arr(2))) .collectAsMap() // 广播这个映射到所有Executor节点 val broadcastMappings = sc.broadcast(mappingsMap)
然后修改_getHiveTypeMapping函数,用广播变量查询映射:
def _getHiveTypeMapping(dataType: String, dataFormat: String) : String = { var definedType="" // 新增数字判断工具函数 def isNumeric(str: String): Boolean = { try { str.toDouble true } catch { case _: NumberFormatException => false } } try { if(dataFormat != null && dataFormat.toUpperCase.contains("DATETIME")){ definedType="datetime" } else if(dataFormat != null && isNumeric(dataFormat)) { definedType="Double" } else { definedType="Unknown" } } catch { case _: Throwable => definedType="Unknown" } // 处理默认情况,空值也要做特殊处理 val finalDefinedType = if (definedType.isEmpty || definedType == "Unknown") { dataFormat match { case null => "" case _ => dataFormat } } else { definedType } // 用广播的映射查询结果 try { val key = (dataType.toUpperCase, finalDefinedType.toUpperCase) broadcastMappings.value.getOrElse(key, "") } catch { case e: Exception => e.getMessage } }
3. 修正数字判断的逻辑
原代码里dataFormat.toDouble.getClass.getName == "double"是错误的:dataFormat.toDouble得到的是java.lang.Double类型,它的类名是"java.lang.Double",不是基本类型的"double"。换成上面的isNumeric函数来判断字符串是否为数字格式,才是正确的做法。
验证效果
修改完这些之后,再执行collect()就不会触发NPE了,而且性能也会因为复用广播的映射而提升不少。
内容的提问来源于stack exchange,提问作者Arun

