从Spark DataFrame中的单个列派生多个列

ssh*_*off 48 scala user-defined-functions dataframe apache-spark apache-spark-sql

我有一个带有巨大可解析元数据的DF作为Dataframe中的单个字符串列,我们称之为DFA,使用ColmnA.

我想打破这一列,将ColmnA分成多个列,通过一个函数,ClassXYZ = Func1(ColmnA).此函数返回一个具有多个变量的类ClassXYZ,现在每个变量都必须映射到新列,例如ColmnA1,ColmnA2等.

如何通过调用此Func1一次,使用这些附加列从一个Dataframe到另一个Data转换,而不必重复它来创建所有列.

如果我每次都要调用这个巨大的函数添加一个新列,它很容易解决,但这是我希望避免的.

请使用工作或伪代码建议.

谢谢

桑杰

zer*_*323 64

一般来说,你想要的东西不是直接可能的.UDF当时只能返回一列.您可以通过两种不同的方式克服此限制:

  1. 返回复杂类型的列.最通用的解决方案是a,StructType但你可以考虑ArrayType或者MapType也可以.

    import org.apache.spark.sql.functions.udf
    
    val df = Seq(
      (1L, 3.0, "a"), (2L, -1.0, "b"), (3L, 0.0, "c")
    ).toDF("x", "y", "z")
    
    case class Foobar(foo: Double, bar: Double)
    
    val foobarUdf = udf((x: Long, y: Double, z: String) => 
      Foobar(x * y, z.head.toInt * y))
    
    val df1 = df.withColumn("foobar", foobarUdf($"x", $"y", $"z"))
    df1.show
    // +---+----+---+------------+
    // |  x|   y|  z|      foobar|
    // +---+----+---+------------+
    // |  1| 3.0|  a| [3.0,291.0]|
    // |  2|-1.0|  b|[-2.0,-98.0]|
    // |  3| 0.0|  c|   [0.0,0.0]|
    // +---+----+---+------------+
    
    df1.printSchema
    // root
    //  |-- x: long (nullable = false)
    //  |-- y: double (nullable = false)
    //  |-- z: string (nullable = true)
    //  |-- foobar: struct (nullable = true)
    //  |    |-- foo: double (nullable = false)
    //  |    |-- bar: double (nullable = false)
    
    Run Code Online (Sandbox Code Playgroud)

    这可以在以后容易地变平,但通常不需要这样做.

  2. 切换到RDD,重塑和重建DF:

    import org.apache.spark.sql.types._
    import org.apache.spark.sql.Row
    
    def foobarFunc(x: Long, y: Double, z: String): Seq[Any] = 
      Seq(x * y, z.head.toInt * y)
    
    val schema = StructType(df.schema.fields ++
      Array(StructField("foo", DoubleType), StructField("bar", DoubleType)))
    
    val rows = df.rdd.map(r => Row.fromSeq(
      r.toSeq ++
      foobarFunc(r.getAs[Long]("x"), r.getAs[Double]("y"), r.getAs[String]("z"))))
    
    val df2 = sqlContext.createDataFrame(rows, schema)
    
    df2.show
    // +---+----+---+----+-----+
    // |  x|   y|  z| foo|  bar|
    // +---+----+---+----+-----+
    // |  1| 3.0|  a| 3.0|291.0|
    // |  2|-1.0|  b|-2.0|-98.0|
    // |  3| 0.0|  c| 0.0|  0.0|
    // +---+----+---+----+-----+
    
    Run Code Online (Sandbox Code Playgroud)

  • 当你说"通常没有[扁平化一列]"时,为什么会这样?或者,spark是否允许您使用顶级列执行的大多数操作也可以使用分层数据(例如`df1.foobar.foo`)完成? (2认同)
  • @max因为简单的`structs`可以在几乎任何上下文中使用,当人们通常使用平面结构时(使用简单的点语法`fooobar.foo`).但它不适用于集合类型.您也可以访问http://stackoverflow.com/a/33850490/1560062 (2认同)

Edw*_*Guo 16

假设在你的函数之后会有一系列元素,给出一个例子如下:

val df = sc.parallelize(List(("Mike,1986,Toronto", 30), ("Andre,1980,Ottawa", 36), ("jill,1989,London", 27))).toDF("infoComb", "age")
df.show
+------------------+---+
|          infoComb|age|
+------------------+---+
|Mike,1986,Toronto| 30|
| Andre,1980,Ottawa| 36|
|  jill,1989,London| 27|
+------------------+---+
Run Code Online (Sandbox Code Playgroud)

现在你可以用这个infoComb做的是你可以开始拆分字符串并获得更多列:

df.select(expr("(split(infoComb, ','))[0]").cast("string").as("name"), expr("(split(infoComb, ','))[1]").cast("integer").as("yearOfBorn"), expr("(split(infoComb, ','))[2]").cast("string").as("city"), $"age").show
+-----+----------+-------+---+
| name|yearOfBorn|   city|age|
+-----+----------+-------+---+
|Mike|      1986|Toronto| 30|
|Andre|      1980| Ottawa| 36|
| jill|      1989| London| 27|
+-----+----------+-------+---+
Run Code Online (Sandbox Code Playgroud)

希望这可以帮助.


Nie*_*and 5

如果生成的列与原始列的长度相同,则可以使用withColumn函数和应用udf创建全新列.在此之后,您可以删除原始列,例如:

 val newDf = myDf.withColumn("newCol1", myFun(myDf("originalColumn")))
.withColumn("newCol2", myFun2(myDf("originalColumn"))
.drop(myDf("originalColumn"))
Run Code Online (Sandbox Code Playgroud)

其中myFun是一个定义如下的udf:

   def myFun= udf(
    (originalColumnContent : String) =>  {
      // do something with your original column content and return a new one
    }
  )
Run Code Online (Sandbox Code Playgroud)