【问题标题】:Change nullable property of column in spark dataframe更改火花数据框中列的可为空属性
【发布时间】:2016-01-16 14:13:14
【问题描述】:

我正在为一些测试手动创建一个数据框。创建它的代码是:

case class input(id:Long, var1:Int, var2:Int, var3:Double)
val inputDF = sqlCtx
  .createDataFrame(List(input(1110,0,1001,-10.00),
    input(1111,1,1001,10.00),
    input(1111,0,1002,10.00)))

所以架构看起来像这样:

root
 |-- id: long (nullable = false)
 |-- var1: integer (nullable = false)
 |-- var2: integer (nullable = false)
 |-- var3: double (nullable = false)

我想为这些变量中的每一个设置“nullable = true”。如何从一开始就声明它或在创建后将其切换到新数据框中?

【问题讨论】:

  • input 应改为 Input - 案例类名应大写(标题大小写)

标签: scala apache-spark spark-dataframe


【解决方案1】:

回答

进口

import org.apache.spark.sql.types.{StructField, StructType}
import org.apache.spark.sql.{DataFrame, SQLContext}
import org.apache.spark.{SparkConf, SparkContext}

你可以使用

/**
 * Set nullable property of column.
 * @param df source DataFrame
 * @param cn is the column name to change
 * @param nullable is the flag to set, such that the column is  either nullable or not
 */
def setNullableStateOfColumn( df: DataFrame, cn: String, nullable: Boolean) : DataFrame = {

  // get schema
  val schema = df.schema
  // modify [[StructField] with name `cn`
  val newSchema = StructType(schema.map {
    case StructField( c, t, _, m) if c.equals(cn) => StructField( c, t, nullable = nullable, m)
    case y: StructField => y
  })
  // apply new schema
  df.sqlContext.createDataFrame( df.rdd, newSchema )
}

直接。

您还可以通过“pimp my library”库模式使该方法可用(请参阅我的 SO 帖子 What is the best way to define custom methods on a DataFrame?),这样您就可以调用

val df = ....
val df2 = df.setNullableStateOfColumn( "id", true )

编辑

替代解决方案 1

使用setNullableStateOfColumn的轻微修改版本

def setNullableStateForAllColumns( df: DataFrame, nullable: Boolean) : DataFrame = {
  // get schema
  val schema = df.schema
  // modify [[StructField] with name `cn`
  val newSchema = StructType(schema.map {
    case StructField( c, t, _, m) ⇒ StructField( c, t, nullable = nullable, m)
  })
  // apply new schema
  df.sqlContext.createDataFrame( df.rdd, newSchema )
}

替代方案 2

明确定义架构。 (使用反射创建更通用的解决方案)

configuredUnitTest("Stackoverflow.") { sparkContext =>

  case class Input(id:Long, var1:Int, var2:Int, var3:Double)

  val sqlContext = new SQLContext(sparkContext)
  import sqlContext.implicits._


  // use this to set the schema explicitly or
  // use refelection on the case class member to construct the schema
  val schema = StructType( Seq (
    StructField( "id", LongType, true),
    StructField( "var1", IntegerType, true),
    StructField( "var2", IntegerType, true),
    StructField( "var3", DoubleType, true)
  ))

  val is: List[Input] = List(
    Input(1110, 0, 1001,-10.00),
    Input(1111, 1, 1001, 10.00),
    Input(1111, 0, 1002, 10.00)
  )

  val rdd: RDD[Input] =  sparkContext.parallelize( is )
  val rowRDD: RDD[Row] = rdd.map( (i: Input) ⇒ Row(i.id, i.var1, i.var2, i.var3))
  val inputDF = sqlContext.createDataFrame( rowRDD, schema ) 

  inputDF.printSchema
  inputDF.show()
}

【讨论】:

  • 所以没有办法只对列进行全面重置?如果需要,我总是可以将列名抓取到一个列表中并循环遍历该列表。顺便说一句,“皮条客我的图书馆”这件事太棒了!
  • 啊,现在我明白你的意思了。您可以通过StructTypecreateDataFrame 指定架构。将对我的答案进行编辑。
  • 所有这些都是为了启用几乎所有 SQL 引擎中的典型默认行为:字段可以包含空值?
  • 我的观察是,这会在源 RDD 上创建一个逻辑计划,从而导致额外的处理 - 这似乎被认为是一个动作,因为我的阶段现在停止在 createDataFrame 行,而不是一些后期处理阶段。
【解决方案2】:

另一种选择,如果您需要就地更改数据框,并且无法重新创建,您可以执行以下操作:

.withColumn("col_name", when(col("col_name").isNotNull, col("col_name")).otherwise(lit(null)))

然后Spark 会认为该列可能包含null,并且可空性将设置为true。 此外,您可以使用udf 将您的值包装在Option 中。 即使对于流媒体案例也能正常工作。

【讨论】:

  • 知道如何在结构化流数据帧中实现逆向(将列设置为不可为空)吗?
  • 不错! PySpark 版本是.withColumn("col_name", when(col("col_name").isNotNull(), col("col_name")).otherwise(lit(None)))
  • otherwise 似乎不需要。它显示在this answer
  • 这应该是公认的答案
【解决方案3】:

这是一个迟到的答案,但想为来这里的人提供一个替代解决方案。通过对代码进行以下修改,您可以从一开始就自动使 DataFrame Column 为空:

case class input(id:Option[Long], var1:Option[Int], var2:Int, var3:Double)
val inputDF = sqlContext
  .createDataFrame(List(input(Some(1110),Some(0),1001,-10.00),
    input(Some(1111),Some(1),1001,10.00),
    input(Some(1111),Some(0),1002,10.00)))
inputDF.printSchema

这将产生:

root
 |-- id: long (nullable = true)
 |-- var1: integer (nullable = true)
 |-- var2: integer (nullable = false)
 |-- var3: double (nullable = false)

defined class input
inputDF: org.apache.spark.sql.DataFrame = [id: bigint, var1: int, var2: int, var3: double]

基本上,如果您通过使用Some([element])None 作为实际输入将字段声明为Option,则该字段可以为空。否则,该字段将不能为空。我希望这会有所帮助!

【讨论】:

    【解决方案4】:

    设置所有列可空参数的更紧凑版本

    可以使用_.copy(nullable = nullable) 代替case StructField( c, t, _, m) ⇒ StructField( c, t, nullable = nullable, m)。那么整个函数可以写成:

    def setNullableStateForAllColumns( df: DataFrame, nullable: Boolean) : DataFrame = {
      df.sqlContext.createDataFrame(df.rdd, StructType(df.schema.map(_.copy(nullable = nullable))))
    }
    

    【讨论】:

      【解决方案5】:

      只需在您的案例类中使用 java.lang.Integer 而不是 scala.Int。

      case class input(id:Long, var1:java.lang.Integer , var2:java.lang.Integer , var3:java.lang.Double)
      

      【讨论】:

      • 先生,为什么我们需要更喜欢 java Integer 而不是 spark 的任何具体原因?
      【解决方案6】:

      谢谢Martin Senne。 只是一点点补充。对于内部结构类型,您可能需要递归设置 nullable,如下所示:

      def setNullableStateForAllColumns(df: DataFrame, nullable: Boolean): DataFrame = {
          def set(st: StructType): StructType = {
            StructType(st.map {
              case StructField(name, dataType, _, metadata) =>
                val newDataType = dataType match {
                  case t: StructType => set(t)
                  case _ => dataType
                }
                StructField(name, newDataType, nullable = nullable, metadata)
            })
          }
      
          df.sqlContext.createDataFrame(df.rdd, set(df.schema))
        }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-11-01
        • 1970-01-01
        • 1970-01-01
        • 2018-11-30
        • 1970-01-01
        • 1970-01-01
        • 2020-04-20
        • 1970-01-01
        相关资源
        最近更新 更多