在Spark-SQL中创建用户定义的函数

use*_*024 18 sql apache-spark

我是新手来激发和激发sql,我试图使用spark SQL查询一些数据.

我需要从作为字符串给出的日期中获取月份.

我认为不可能直接从sparkqsl查询月份,所以我想在scala中编写用户定义的函数.

是否有可能在sparkSQL中编写udf,如果可能,任何人都可以提出编写udf的最佳方法.

请帮忙

Spi*_*lov 11

如果您愿意使用语言集成查询,则可以执行此操作,至少是为了过滤.

对于包含以下内容的数据文件dates.txt:

one,2014-06-01
two,2014-07-01
three,2014-08-01
four,2014-08-15
five,2014-09-15
Run Code Online (Sandbox Code Playgroud)

您可以根据需要在UDF中打包尽可能多的Scala日期魔法,但我会保持简单:

def myDateFilter(date: String) = date contains "-08-"
Run Code Online (Sandbox Code Playgroud)

将其全部设置如下 - 其中很多内容来自编程指南.

val sqlContext = new org.apache.spark.sql.SQLContext(sc)
import sqlContext._

// case class for your records
case class Entry(name: String, when: String)

// read and parse the data
val entries = sc.textFile("dates.txt").map(_.split(",")).map(e => Entry(e(0),e(1)))
Run Code Online (Sandbox Code Playgroud)

您可以将UDF用作WHERE子句的一部分:

val augustEntries = entries.where('when)(myDateFilter).select('name, 'when)
Run Code Online (Sandbox Code Playgroud)

并看到结果:

augustEntries.map(r => r(0)).collect().foreach(println)
Run Code Online (Sandbox Code Playgroud)

请注意where我使用的方法的版本,在doc中声明如下:

def where[T1](arg1: Symbol)(udf: (T1) ? Boolean): SchemaRDD
Run Code Online (Sandbox Code Playgroud)

因此,UDF只能接受一个参数,但您可以组合多个.where()调用来过滤多个列.

编辑Spark 1.2.0(真的是1.1.0)

虽然它没有真正记录,但Spark现在支持注册UDF,因此可以从SQL查询它.

上述UDF可以使用以下方式注册:

sqlContext.registerFunction("myDateFilter", myDateFilter)
Run Code Online (Sandbox Code Playgroud)

如果表已注册

sqlContext.registerRDDAsTable(entries, "entries")
Run Code Online (Sandbox Code Playgroud)

可以使用查询

sqlContext.sql("SELECT * FROM entries WHERE myDateFilter(when)")
Run Code Online (Sandbox Code Playgroud)

有关详细信息,请参阅此示例.