加载Scala样例类时触发NullPointerException问题排查与优化
解决方案
1. 读取后先清洗无效键值对
直接在Spark DataFrame层面清洗fun_temporal_unit字段,过滤掉key或value为null的条目,再转换为目标case class:
import org.apache.spark.sql.functions._ // 假设原始读取的BigQuery数据为df val cleanedDf = df.withColumn( "XPI_MIN", struct( // 保留XPI_MIN的其他字段,示例用other_fields代替 col("XPI_MIN.other_fields"), struct( // 保留functions的其他字段 col("XPI_MIN.functions.other_func_fields"), // 过滤key/value为null的条目,再转成Map格式 expr(""" map_from_entries( filter( transform(XPI_MIN.functions.fun_temporal_unit, e -> struct(e.key as k, e.value as v) ), entry -> entry.k is not null and entry.v is not null ) ) """).alias("fun_temporal_unit") ).alias("functions") ) ) // 原case class定义不变 case class Functions(other_func_fields: String, fun_temporal_unit: Option[Map[String,String]]) case class XPIMin(other_fields: String, functions: Functions) case class YourBizModel(XPI_MIN: Option[XPIMin], /* 其他字段 */) // 转换为case class数据集 val resultDs = cleanedDf.as[YourBizModel]
2. 修改Case Class兼容无效数据
如果不想提前清洗,可先将fun_temporal_unit定义为能接收null的结构,后续在业务逻辑中处理:
// 定义临时接收结构,允许key/value为null case class TemporalEntry(key: Option[String], value: Option[String]) case class Functions(other_func_fields: String, fun_temporal_unit: Option[List[TemporalEntry]]) case class XPIMin(other_fields: String, functions: Functions) case class YourBizModel(XPI_MIN: Option[XPIMin], /* 其他字段 */) // 读取数据后在业务逻辑中过滤无效条目 val cleanedResult = resultDs.map { model => model.XPI_MIN.map { xpiMin => val validTemporalMap = xpiMin.functions.fun_temporal_unit.map { entries => entries.collect { case TemporalEntry(Some(k), Some(v)) => k -> v }.toMap } xpiMin.copy(functions = xpiMin.functions.copy(fun_temporal_unit = validTemporalMap)) }.map(cleanedXpi => model.copy(XPI_MIN = Some(cleanedXpi))).getOrElse(model) }
3. 自定义UDF处理字段转换
编写UDF统一处理包含null的键值对集合,输出合法的Map:
import org.apache.spark.sql.api.java.UDF1 import scala.collection.JavaConverters._ // 自定义UDF:过滤key/value为null的条目,返回Option[Map] val cleanTemporalMapUdf = udf((rawEntries: java.util.List[java.util.Map[String, String]]) => { Option(rawEntries).map { entryList => entryList.asScala.flatMap { entry => val key = entry.get("key") val value = entry.get("value") if (key != null && value != null) Some(key -> value) else None }.toMap } }) // 应用UDF清洗字段 val cleanedDf = df.withColumn( "XPI_MIN", struct( col("XPI_MIN.*"), struct( col("XPI_MIN.functions.*"), cleanTemporalMapUdf(col("XPI_MIN.functions.fun_temporal_unit")).alias("fun_temporal_unit") ).alias("functions") ) ) // 转换为原case class val resultDs = cleanedDf.as[YourBizModel]
内容的提问来源于stack exchange,提问作者Vikrant Singh Rana
相关产品推荐
相关产品推荐

