【发布时间】:2016-12-24 16:26:28
【问题描述】:
使用 Spark 的 DataFrame 时,需要使用用户定义函数 (UDF) 来映射列中的数据。 UDF 要求显式指定参数类型。就我而言,我需要操作由对象数组组成的列,但我不知道要使用什么类型。这是一个例子:
import sqlContext.implicits._
// Start with some data. Each row (here, there's only one row)
// is a topic and a bunch of subjects
val data = sqlContext.read.json(sc.parallelize(Seq(
"""
|{
| "topic" : "pets",
| "subjects" : [
| {"type" : "cat", "score" : 10},
| {"type" : "dog", "score" : 1}
| ]
|}
""")))
使用内置的org.apache.spark.sql.functions对列中的数据进行基本操作比较简单
import org.apache.spark.sql.functions.size
data.select($"topic", size($"subjects")).show
+-----+--------------+
|topic|size(subjects)|
+-----+--------------+
| pets| 2|
+-----+--------------+
编写自定义 UDF 来执行任意操作通常很容易
import org.apache.spark.sql.functions.udf
val enhance = udf { topic : String => topic.toUpperCase() }
data.select(enhance($"topic"), size($"subjects")).show
+----------+--------------+
|UDF(topic)|size(subjects)|
+----------+--------------+
| PETS| 2|
+----------+--------------+
但是,如果我想使用 UDF 来操作“主题”列中的对象数组怎么办?我对 UDF 中的参数使用什么类型?比如我想重新实现size函数,而不是使用spark提供的那个:
val my_size = udf { subjects: Array[Something] => subjects.size }
data.select($"topic", my_size($"subjects")).show
显然Array[Something] 不起作用...我应该使用什么类型!?我应该完全放弃Array[] 吗?闲逛告诉我scala.collection.mutable.WrappedArray 可能与它有关,但我仍然需要提供另一种类型。
【问题讨论】:
标签: scala apache-spark dataframe apache-spark-sql user-defined-functions