【问题标题】:How to separate array or vector column into multiple columns?如何将数组或向量列分成多列?
【发布时间】:2017-10-05 02:01:39
【问题描述】:

假设我有一个 Spark Dataframe 生成为:

val df = Seq(
    (Array(1, 2, 3), Array("a", "b", "c")),
    (Array(1, 2, 3), Array("a", "b", "c"))
).toDF("Col1", "Col2")

可以在"Col1" 的第一个索引处提取元素,例如:

val extractFirstInt = udf { (x: Seq[Int], i: Int) => x(i) }
df.withColumn("Col1_1", extractFirstInt($"Col1", lit(1)))

对于第二列 "Col2" 也类似,例如

val extractFirstString = udf { (x: Seq[String], i: Int) => x(i) }
df.withColumn("Col2_1", extractFirstString($"Col2", lit(1)))

但是代码重复有点难看——我需要为每个底层元素类型单独的 UDF。

有没有办法编写一个 generic UDF,自动推断 Spark 数据集列中底层数组的类型?例如。我希望能够编写类似 (pseudocode; with generic T)

val extractFirst = udf { (x: Seq[T], i: Int) => x(i) }
df.withColumn("Col1_1", extractFirst($"Col1", lit(1)))

Spark / Scala 编译器会以某种方式自动推断出 T 类型(如果合适,可能会使用反射)。

如果您知道一个解决方案既适用于数组列,也适用于 Spark 自己的 DenseVector / SparseVector 类型,则可以加分。我想避免的主要事情(如果可能的话)是为我要处理的每个底层数组元素类型定义一个单独的 UDF。

【问题讨论】:

  • 为什么要使用 udf?对于向量,您没有任何合理的选择,但数组支持 getItemapply...
  • 好点!也许最好的解决方案是检测列是否包含ArrayType 对象,如果是,则使用.getItem() 和UDF(用于向量)。
  • 是的,而且几乎是免费的。由于 Vector 不是原生类型,因此比较棘手,但从好的方面来说,您只需担心一种类型。

标签: scala apache-spark


【解决方案1】:

也许frameless 可能是一个解决方案?

由于操作数据集需要给定类型的Encoder,因此您必须预先定义类型,以便 Spark SQL 可以为您创建一个。我认为生成各种编码器支持的类型的 Scala 宏在这里是有意义的。

截至目前,我会为每个类型定义一个泛型方法和一个 UDF(这与您希望找到一种方法让 " 一个泛型 UDF 相违背,它会自动推断底层 Array 的类型Spark 数据集的列”)。

def myExtract[T](x: Seq[T], i: Int) = x(i)
// define UDF for extracting strings
val extractString = udf(myExtract[String] _)

如下使用:

val df = Seq(
    (Array(1, 2, 3), Array("a", "b", "c")),
    (Array(1, 2, 3), Array("a", "b", "c"))
).toDF("Col1", "Col2")

scala> df.withColumn("Col1_1", extractString($"Col2", lit(1))).show
+---------+---------+------+
|     Col1|     Col2|Col1_1|
+---------+---------+------+
|[1, 2, 3]|[a, b, c]|     b|
|[1, 2, 3]|[a, b, c]|     b|
+---------+---------+------+

您可以改为探索Dataset(不是DataFrame,即Dataset[Row])。这将为您提供所有类型机器(也许您可以避免任何宏开发)。

【讨论】:

  • 数组中的元素数量是否可能未知?例如,如果在 col2 中有一个大小未知的数组。
【解决方案2】:

根据@zero323 的建议,我专注于以下形式的实现:

def extractFirst(df: DataFrame, column: String, into: String) = {

  // extract column of interest
  val col = df.apply(column)

  // figure out the type name for this column
  val schema = df.schema
  val typeName = schema.apply(schema.fieldIndex(column)).dataType.typeName

  // delegate based on column type
  typeName match {

    case "array"  => df.withColumn(into, col.getItem(0))
    case "vector" => {
      // construct a udf to extract first element
      // (could almost certainly do better here,
      // but this demonstrates the strategy regardless)
      val extractor = udf {
        (x: Any) => {
          val el = x.getClass.getDeclaredMethod("toArray").invoke(x)
          val array = el.asInstanceOf[Array[Double]]
          array(0)
        }
      }

      df.withColumn(into, extractor(col))
    }

    case _ => throw new IllegalArgumentException("unexpected type '" + typeName + "'")
  }
}

【讨论】:

    猜你喜欢
    • 2016-09-15
    • 2021-11-01
    • 1970-01-01
    • 1970-01-01
    • 2014-02-25
    • 2020-06-15
    • 1970-01-01
    • 2020-09-03
    • 2017-05-03
    相关资源
    最近更新 更多