【发布时间】: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?对于向量,您没有任何合理的选择,但数组支持
getItem和apply... -
好点!也许最好的解决方案是检测列是否包含
ArrayType对象,如果是,则使用.getItem()和UDF(用于向量)。 -
是的,而且几乎是免费的。由于 Vector 不是原生类型,因此比较棘手,但从好的方面来说,您只需担心一种类型。
标签: scala apache-spark