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

Spark DataFrame中from_json解析Int字段为null的问题排查

Spark中from_json解析Int类型字段为null的问题

问题现象

在Spark DataFrame中通过Encoders.product从case class生成Schema,再使用from_json解析JSON字符串时出现以下异常:

  • 当case class中PRICE字段定义为Option[Int]时,解析后该字段值为null;
  • 当PRICE定义为Option[String]时,可正常获取字符串类型的字段值;
  • 通过Encoders.product生成的Int类型Schema结构是正确的。

核心代码片段

case class MyProducts(PRODUCT_ID: Option[String], DESCRIPTION: Option[String], PRICE: Option[Int], OLD_FIELD_1: Option[String]) 
val ProductsSchema = Encoders.product[MyProducts].schema

val df_products_output_final = df_products_output.withColumn("parsedProducts", from_json(col("afterImage"), ProductsSchema)) 

完整代码

import org.json4s._
import org.json4s.jackson.JsonMethods._
import spark.implicits._
import org.apache.spark.sql.functions.{col, lit, when, from_json, map_keys, map_values, regexp_replace, coalesce}
import org.apache.spark.sql.types.{MapType, StringType}
import org.apache.spark.sql.Encoders
import org.apache.spark.sql.types.{MapType, StringType, StructType, IntegerType}

case class MyMeta(op: String, table: String)
val metaSchema = Encoders.product[MyMeta].schema
case class MySales(NUM: Option[Integer], PRODUCT_ID: Option[String], DESCRIPTION: Option[String], OLD_FIELD_1: Option[String]) 
val salesSchema = Encoders.product[MySales].schema
case class MyProducts(PRODUCT_ID: Option[String], DESCRIPTION: Option[String], PRICE: Option[Int], OLD_FIELD_1: Option[String]) 
val ProductsSchema = Encoders.product[MyProducts].schema

def getAfterImage (op: String, data: String, key: String, jsonOLD_TABLE_FIELDS: String) : String = {   
  val jsonOLD_FIELDS = parse(jsonOLD_TABLE_FIELDS)   
  val jsonData = parse(data)                          
  val jsonKey = parse(key)                           
   
  op match {
  case "ins" =>
               return(compact(render(jsonData merge jsonOLD_FIELDS)))
  case _ => 
               val Diff(changed, added, deleted) = jsonKey diff jsonData
               return(compact(render(changed merge deleted merge jsonOLD_FIELDS)))
  }
}
val afterImage = spark.udf.register("callUDFAI", getAfterImage _)

val path = "/FileStore/tables/json_0006_file.txt"  
val df = spark.read.text(path)  // String.
val df2 = df.withColumn("value", from_json(col("value"), MapType(StringType, StringType)))    
val df3 = df2.select(map_values(col("value")))  
val df4 = df3.select($"map_values(value)"(0).as("meta"), $"map_values(value)"(1).as("data"), $"map_values(value)"(2).as("key")).withColumn("parsedMeta", from_json(col("meta"), metaSchema)).drop("meta").select(col("parsedMeta.*"), col("data"), col("key")).withColumn("key2", coalesce(col("key"), lit(""" { "DUMMY_FIELD_XXX": ""} """))).toDF().cache()

val df_sales    = df4.filter('table === "BILL.SALES") 
val df_products = df4.filter('table === "BILL.PRODUCTS")
val df_sales_output = df_sales.withColumn("afterImage", afterImage(col("op"), col("data"), col("key2") , lit(""" { "OLD_FIELD_1": ""} """)))
                              .select("afterImage") 
val df_products_output = df_products.withColumn("afterImage", afterImage(col("op"), col("data"), col("key2") , lit(""" { "OLD_FIELD_A":"", "OLD_FIELD_B":""} """)))
                                    .select("afterImage")                          
val df_sales_output_final = df_sales_output.withColumn("parsedSales", from_json(col("afterImage"), salesSchema)) 
df_products_output_final.show(false)
df_products_output_final.printSchema()

问题原因

  1. JSON数据类型不匹配:afterImage字段中的PRICE值是字符串格式的数字(例如"PRICE": "99"),而非原始的数字类型("PRICE": 99)。Spark的from_json函数类型校验严格,当Schema定义为IntegerType时,无法自动将字符串类型的数字转换为Int,解析失败后返回null;而定义为StringType时,字符串格式的数字可直接匹配,解析正常。
  2. UDF处理导致类型转换:自定义UDFgetAfterImage使用json4s处理JSON时,可能将原始的数字类型值转换为了字符串类型(比如在merge操作过程中,将数字转成了带引号的字符串),进一步加剧了类型不匹配问题。

解决建议

  • 修正UDF逻辑:调整json4s的处理方式,确保PRICE字段在生成的JSON中保持原始数字类型;
  • 显式类型转换:先将PRICE解析为String类型,再通过cast(IntegerType)转换为Int;
  • 宽松解析配置:Spark 3.0+可通过from_json的options参数开启类型自动转换(例如from_json(col("afterImage"), ProductsSchema, Map("mode" -> "PERMISSIVE"))),但更推荐显式处理数据类型以避免潜在问题。

内容的提问来源于stack exchange,提问作者Ged

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:27:21