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

如何创建含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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 09:33:20