Spark Scala构建含Map嵌套Struct的测试DataFrame遇空值问题
问题分析与解决方案
问题根源
你构造的JSON字符串不符合标准JSON语法,里面使用了Spark的Row()对象初始化语法(比如Row(0.7, 1.3)),但from_json函数只能解析标准JSON结构,无法识别这种Scala代码语法,因此解析失败导致scores列全为null。
解决方案
针对单元测试场景,推荐两种可靠的实现方式:
方式一:直接构造符合Schema的DataFrame(最适合单元测试)
不需要通过JSON中转,直接用Spark API构造匹配目标Schema的数据,避免解析错误:
import org.apache.spark.sql.types._ import org.apache.spark.sql.SparkSession // 初始化SparkSession(单元测试中通常已注入) val spark = SparkSession.builder().master("local").getOrCreate() import spark.implicits._ // 定义目标Schema val targetSchema = StructType(Seq( StructField("ids", StringType), StructField("scores", MapType( StringType, StructType(Seq( StructField("myscore1", DoubleType), StructField("myscore2", DoubleType) )) )) )) // 直接构造数据 val exDf = spark.createDataFrame(Seq( ( "id_1", Map( "key1" -> (0.7, 1.3), "key2" -> (0.5, 1.2) ) ) )).toDF("ids", "scores") // 验证结果 exDf.printSchema() exDf.show(false)
方式二:修正JSON字符串后用from_json解析
如果必须通过JSON字符串构造,需要将字符串改为标准JSON格式(用键值对表示Struct结构):
import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().master("local").getOrCreate() import spark.implicits._ // 定义Struct和Map的Schema val scoreStructSchema = StructType(Seq( StructField("myscore1", DoubleType), StructField("myscore2", DoubleType) )) val scoreMapSchema = MapType(StringType, scoreStructSchema) // 使用标准JSON字符串 val exDf = List( ("id_1", """{"key1":{"myscore1":0.7,"myscore2":1.3}, "key2":{"myscore1":0.5,"myscore2":1.2}}""") ).toDF("ids", "scores_str") .withColumn("scores", from_json(col("scores_str"), scoreMapSchema)) .drop("scores_str") // 验证结果 exDf.printSchema() exDf.show(false)
验证结果
两种方式生成的DataFrame Schema与目标完全匹配,数据输出如下:
+-----+---------------------------------------------------+ |ids |scores | +-----+---------------------------------------------------+ |id_1 |{key1 -> {0.7, 1.3}, key2 -> {0.5, 1.2}} | +-----+---------------------------------------------------+
内容的提问来源于stack exchange,提问作者Callie
相关产品推荐
相关产品推荐

