Spark DataFrame Schema不匹配:嵌套Struct字段空值引发运行时异常
Spark DataFrame Schema转换非空约束异常解决
原DataFrame数据
+----------+-----------------------+----------------------+-----------+------+-------------+------------+ |TECHNOLOGY|KPI_NAME |FUNCTIONS |DESCRIPTION|ACTION|FORMULA_VALID|VALIDITY_LOG| +----------+-----------------------+----------------------+-----------+------+-------------+------------+ |GSM |Cell_Availability_test3|{SUM, SUM, NULL, NULL}|NULL |ADD |true |[] | +----------+-----------------------+----------------------+-----------+------+-------------+------------+
原DataFrame Schema
root |-- TECHNOLOGY: string (nullable = true) |-- KPI_NAME: string (nullable = true) |-- FUNCTIONS: struct (nullable = true) | |-- fun_temporal: string (nullable = true) | |-- fun_regional: string (nullable = true) | |-- fun_temporal_unit: map (nullable = true) | | |-- key: string | | |-- value: string (valueContainsNull = true) | |-- fun_regional_unit: map (nullable = true) | | |-- key: string | | |-- value: string (valueContainsNull = true) |-- DESCRIPTION: string (nullable = true) |-- ACTION: string (nullable = true) |-- FORMULA_VALID: boolean (nullable = false) |-- VALIDITY_LOG: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- key: string (nullable = true) | | |-- value: string (nullable = true)
目标Schema定义及转换代码
val outputTypeTest: StructType = StructType(Seq( StructField("TECHNOLOGY", StringType, true), StructField("KPI_NAME", StringType, true), StructField("FUNCTIONS", StructType(Seq( StructField("fun_temporal", StringType, true), StructField("fun_regional", StringType, true), StructField("fun_temporal_unit", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), false), StructField("fun_regional_unit", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), false))), true), StructField("DESCRIPTION", StringType, true), StructField("ACTION", StringType, true), StructField("FORMULA_VALID", BooleanType, true), StructField("VALIDITY_LOG", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), false))) val formulaMappingOutputNotTypedTest= formulaMappingOutputNotTyped.select("TECHNOLOGY","KPI_NAME","FUNCTIONS","DESCRIPTION","ACTION","FORMULA_VALID","VALIDITY_LOG") formulaMappingOutputNotTypedTest.show(truncate = false) val formulaMappingOutput = spark.createDataFrame(formulaMappingOutputNotTypedTest.rdd, outputTypeTest)
触发的异常信息
导致错误的原因:java.lang.RuntimeException: 输入行的第2个字段'fun_temporal_unit'不能为null。
问题根源
- 目标Schema中
fun_temporal_unit和fun_regional_unit的nullable参数被设为false,但原DataFrame里这两个字段的值是NULL,直接违反了非空约束。 - 原字段类型是
Map,目标类型是Array[Struct],直接用spark.createDataFrame(rdd, schema)无法自动完成类型转换,还会触发严格的非空校验。
解决步骤
- 调整nullable属性:要么将目标Schema中
fun_temporal_unit和fun_regional_unit的nullable改为true,允许空值;要么提前将原数据中的NULL替换为符合要求的空数组。 - 类型转换:使用Spark内置函数将Map类型字段转换为目标的数组结构体类型,不能直接强制Schema。
修正后的代码示例
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 修正目标Schema:允许fun_temporal_unit和fun_regional_unit为空 val outputTypeTest: StructType = StructType(Seq( StructField("TECHNOLOGY", StringType, true), StructField("KPI_NAME", StringType, true), StructField("FUNCTIONS", StructType(Seq( StructField("fun_temporal", StringType, true), StructField("fun_regional", StringType, true), StructField("fun_temporal_unit", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), true), // nullable改为true StructField("fun_regional_unit", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), true) // nullable改为true )), true), StructField("DESCRIPTION", StringType, true), StructField("ACTION", StringType, true), StructField("FORMULA_VALID", BooleanType, true), StructField("VALIDITY_LOG", ArrayType(StructType(Seq( StructField("key", StringType, true), StructField("value", StringType, true))), false), false) )) // 转换FUNCTIONS中的Map字段为Array[Struct],同时处理NULL值 val transformedDF = formulaMappingOutputNotTypedTest.withColumn("FUNCTIONS", struct( col("FUNCTIONS.fun_temporal"), col("FUNCTIONS.fun_regional"), // Map转Array[Struct],NULL则替换为空数组 when(col("FUNCTIONS.fun_temporal_unit").isNull, array()).otherwise(map_entries(col("FUNCTIONS.fun_temporal_unit"))).alias("fun_temporal_unit"), when(col("FUNCTIONS.fun_regional_unit").isNull, array()).otherwise(map_entries(col("FUNCTIONS.fun_regional_unit"))).alias("fun_regional_unit") )) // 应用目标Schema生成新DataFrame val formulaMappingOutput = spark.createDataFrame(transformedDF.rdd, outputTypeTest)
补充说明
map_entries函数可直接将Map类型转为Array[Struct(key, value)],完美匹配目标Schema的类型要求。- 如果业务规则不允许
fun_temporal_unit为空,无需修改nullable属性,只需确保转换时将NULL替换为业务认可的默认值(比如空数组)即可。
内容的提问来源于stack exchange,提问作者Suhani Bhatia
相关产品推荐
相关产品推荐

