Spark 2.3指定JSON Schema失败:非空列变为可空问题问询
解决Spark 2.3中指定JSON Schema非空列变为可空的问题
我之前也踩过这个坑,Spark 2.x的JSON Reader确实有个“认死理”的特性——哪怕你明确把字段标记为非空(nullable=false),只要读取过程中发现有缺失该字段的行或者字段值为null,它就会自动把该字段的nullable改成true,完全忽略你一开始的设置。结合你的代码,咱们一步步来解决:
第一步:先排查数据本身的问题
首先确认你的JSON文件里是不是真的存在id_str字段缺失或者值为null的情况。可以先临时不加指定Schema读取,看看数据情况:
spark.read.option("multiline", "true").json("/path/to/your/file").select("id_str").show()
如果发现有null或者空值,那问题根源就在数据上——Spark是在帮你规避后续操作因为非空字段出现null而报错的风险。这时候你可以选这几种处理方式:
- 清理脏数据:删除缺失
id_str或者值为null的行; - 接受可空状态:如果业务允许该字段为空,直接保留现状即可;
- 读取后过滤:用
filter(col("id_str").isNotNull)把null值行过滤掉。
第二步:强制修正Schema(如果数据确实无null)
要是确认数据里所有行的id_str都存在且不为null,但Schema还是被改成了可空,那可以在读取后手动修正Schema:
// 先过滤掉可能的null值(保险操作) val filteredDf = mdf.filter(col("id_str").isNotNull) // 重新构建Schema,把id_str的nullable强制设为false val fixedSchema = StructType(filteredDf.schema.map { case StructField(name, dtype, _, meta) if name == "id_str" => StructField(name, dtype, nullable = false, meta) case field => field }) // 用新Schema创建最终的DataFrame val finalDf = spark.createDataFrame(filteredDf.rdd, fixedSchema) // 验证结果 finalDf.printTreeString()
第三步:用Dataset替代DataFrame(更严格的校验)
如果你的数据质量很高,完全不想容忍任何缺失id_str的行,可以用Scala的Case Class定义强类型Schema,通过Dataset的校验机制来保证非空:
// 定义Case Class,字段不加Option,代表该字段必须非空 case class Twitter(id_str: String) // 读取JSON并转换成Dataset val mds = spark.read.option("multiline", "true").json("/path/to/your/file").as[Twitter]
这种方式下,如果JSON里有缺失id_str或者值为null的行,Spark会直接抛出解析异常,强制你处理脏数据,最终得到的Dataset里id_str肯定是非空的。
另外提一句:Spark 2.3是比较老的版本了,这类Schema自动变更的问题在后续的Spark 3.x版本里已经有优化,如果有升级的可能,升级后这类问题会少很多。
内容的提问来源于stack exchange,提问作者Jill Clover
相关产品推荐
相关产品推荐

