如何创建含NaN、null、None值的Spark Scala DataFrame?
问题描述
我在Coursera的Spark+Scala课程中学习了fillna和replace函数的用法,尝试复现验证实际运行效果,但创建包含待替换值的DataFrame时遇到问题。分别尝试使用JSON输入文件和元组序列的方式,均抛出异常。需要指导如何创建同时包含null/NaN/None的DataFrame。
错误代码示例
object HowToCreateDfWithNullsOrNaNs { def main(args: Array[String]): Unit = { fromFile() } def fromFile(): Unit = { // input_file.json: { "name": "Tom", "surname": null, "age": 10} val rddFromJson: RDD[String] = spark.sparkContext.textFile("src/main/resources/input_file.json") import spark.implicits._ /* Exception in thread "main" java.lang.IllegalArgumentException: requirement failed: The number of columns doesn't match. Old column names (1): value New column names (3): name, surname, age */ rddFromJson.toDF("name", "surname", "age") } def fromSeq() = { val tupleSeq: Seq[(String, Any, Int)] = Seq(("Tom", null , 10)) val rdd = spark.sparkContext.parallelize(tupleSeq) /* Exception in thread "main" java.lang.ClassNotFoundException: scala.Any at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:583) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:178) */ import spark.implicits._ rdd.toDF("name", "surname", "age") } }
解决方案
1. 从JSON文件创建包含null的DataFrame
错误原因
用textFile读取JSON得到的是每行一个字符串的RDD,直接调用toDF只会生成一列(默认列名value)的DataFrame,无法直接指定多列名。Spark有专门的JSON解析API来处理这类场景。
修复代码
def fromFileFixed(): Unit = { // 确保input_file.json每行是独立的JSON对象,示例内容:{"name": "Tom", "surname": null, "age": 10} val df = spark.read.json("src/main/resources/input_file.json") df.show() }
2. 从序列创建包含null/NaN/None的DataFrame
错误原因
使用Any作为元组元素,Spark无法序列化该类型,也无法正确推断列类型。三种空值有各自的适用场景:
null:用于String等引用类型NaN:仅适用于Double/Float等浮点类型None:用于Option类型,会被Spark转换为对应基础类型的可空列
修复并扩展代码(包含三种空值)
def fromSeqFixed(): Unit = { import spark.implicits._ // 构造包含三种空值的数据集: // - null:String类型空值 // - Double.NaN:浮点类型非数值 // - None:Option[Int]类型空值,对应Spark的IntegerType可空列 val data = Seq( ("Tom", null, 10.5, None), ("Jerry", "Smith", Double.NaN, Some(15)), ("Alice", "Brown", 20.0, None) ) val df = data.toDF("name", "surname", "weight", "grade") df.show() df.printSchema() }
输出结果
+-----+-------+------+-----+ | name|surname|weight|grade| +-----+-------+------+-----+ | Tom| null| 10.5| null| |Jerry| Smith| NaN| 15| |Alice| Brown| 20.0| null| +-----+-------+------+-----+ root |-- name: string (nullable = true) |-- surname: string (nullable = true) |-- weight: double (nullable = false) |-- grade: integer (nullable = true)
内容的提问来源于stack exchange,提问作者bridgemnc
相关产品推荐
相关产品推荐

