【问题标题】:Why do columns change to nullable in Apache Spark SQL?为什么 Apache Spark SQL 中的列更改为可为空?
【发布时间】:2017-03-28 23:26:03
【问题描述】:

为什么在执行某些函数后使用nullable = true,即使DataFrame中没有NaN值。

val myDf = Seq((2,"A"),(2,"B"),(1,"C"))
         .toDF("foo","bar")
         .withColumn("foo", 'foo.cast("Int"))

myDf.withColumn("foo_2", when($"foo" === 2 , 1).otherwise(0)).select("foo", "foo_2").show

当现在调用df.printSchema 时,nullable 将是false 对于两列。

val foo: (Int => String) = (t: Int) => {
    fooMap.get(t) match {
      case Some(tt) => tt
      case None => "notFound"
    }
  }

val fooMap = Map(
    1 -> "small",
    2 -> "big"
 )
val fooUDF = udf(foo)

myDf
    .withColumn("foo", fooUDF(col("foo")))
    .withColumn("foo_2", when($"foo" === 2 , 1).otherwise(0)).select("foo", "foo_2")
    .select("foo", "foo_2")
    .printSchema

但是现在,nullable 至少有一列是 true,而之前是 false。这怎么解释?

【问题讨论】:

    标签: apache-spark apache-spark-sql apache-spark-dataset


    【解决方案1】:

    您也可以非常快速地更改数据框的架构。像这样的东西可以完成这项工作 -

    def setNullableStateForAllColumns( df: DataFrame, columnMap: Map[String, Boolean]) : DataFrame = {
        import org.apache.spark.sql.types.{StructField, StructType}
        // get schema
        val schema = df.schema
        val newSchema = StructType(schema.map {
        case StructField( c, d, n, m) =>
          StructField( c, d, columnMap.getOrElse(c, default = n), m)
        })
        // apply new schema
        df.sqlContext.createDataFrame( df.rdd, newSchema )
    }
    

    【讨论】:

      【解决方案2】:

      从静态类型结构创建Dataset 时(不依赖于schema 参数)Spark 使用一组相对简单的规则来确定nullable 属性。

      • 如果给定类型的对象可以是null,则其DataFrame 表示为nullable
      • 如果对象是Option[_],那么它的DataFrame 表示是nullable,而None 被认为是SQL NULL
      • 在任何其他情况下,它将被标记为不是nullable

      由于 Scala Stringjava.lang.String,可以是 null,所以生成的列 can 是 nullable。出于同样的原因,bar 列在初始数据集中是nullable

      val data1 = Seq[(Int, String)]((2, "A"), (2, "B"), (1, "C"))
      val df1 = data1.toDF("foo", "bar")
      df1.schema("bar").nullable
      
      Boolean = true
      

      foo 不是(scala.Int 不能是null)。

      df1.schema("foo").nullable
      
      Boolean = false
      

      如果我们将数据定义更改为:

      val data2 = Seq[(Integer, String)]((2, "A"), (2, "B"), (1, "C"))
      

      foo 将是 nullableIntegerjava.lang.Integer 并且装箱的整数可以是 null):

      data2.toDF("foo", "bar").schema("foo").nullable
      
      Boolean = true
      

      另请参阅:SPARK-20668 修改 ScalaUDF 以处理可空性

      【讨论】:

        猜你喜欢
        • 2017-02-05
        • 2011-04-03
        • 1970-01-01
        • 1970-01-01
        • 2016-07-23
        • 1970-01-01
        • 1970-01-01
        • 2013-09-10
        • 2020-10-25
        相关资源
        最近更新 更多