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

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。


问题根源

  1. 目标Schema中fun_temporal_unit和fun_regional_unit的nullable参数被设为false,但原DataFrame里这两个字段的值是NULL,直接违反了非空约束。
  2. 原字段类型是Map,目标类型是Array[Struct],直接用spark.createDataFrame(rdd, schema)无法自动完成类型转换,还会触发严格的非空校验。

解决步骤

  1. 调整nullable属性:要么将目标Schema中fun_temporal_unit和fun_regional_unit的nullable改为true,允许空值;要么提前将原数据中的NULL替换为符合要求的空数组。
  2. 类型转换:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 14:30:55