首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >如何将udf添加到sqlContext中

如何将udf添加到sqlContext中
EN

Stack Overflow用户
提问于 2018-04-13 16:56:31
回答 1查看 1.5K关注 0票数 0

我知道我可以注册一个UDFand函数,因为它可以在SQL查询中使用:

代码语言:javascript
复制
def example(s):
    return len(s)
sqlContext.udf.register("example_udf", example)
spark.sql("SELECT example_udf(col) FROM data")

或者我可以用udf包装Python函数,这样就可以将它应用于dataframe:

代码语言:javascript
复制
from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType
example_udf = udf(example)
data.select(example_udf('col'))

在我的例子中,由于我需要向UDF传递一些其他参数,所以我为UDF构建了一个嵌套函数:

代码语言:javascript
复制
from pyspark.sql.types import BooleanType
from pyspark.sql.functions import col
def my_udf(other_par)
    def example(s):
        return len(s) == other_par
    return udf(example, BooleanType())

dataframe.select(...).where(my_udf(5)(col('col')))

现在我已经有了一个UDF,并且我可以将它应用到dataframe上。但是我也想在spark.sql中使用它,比如第一个块中的SQL查询,而不是dataframe的select或where方法。所以我想知道我是怎么做到的。看起来sqlContext.udf.register只能接受Python函数,而不能接受UDF。

EN

回答 1

Stack Overflow用户

回答已采纳

发布于 2018-04-13 17:08:31

如果您使用的是最新的和最大的(2.3),就不要直接使用udf

代码语言:javascript
复制
def my_udf(other_par, spark):
    def _(s):
        return len(s) == other_par
    return spark.udf.register("my_udf_{}".format(other_par), _, BooleanType())

my_udf_42 = my_udf(42, spark)

spark.sql("SELECT my_udf_42(array(1, 2))").show()
# +----------------------+
# |my_udf_42(array(1, 2))|
# +----------------------+
# |                 false|
# +----------------------+

spark.createDataFrame([([1] * 42, )], ("id", )).select(my_udf_42("id")).show()
# +-------------+
# |my_udf_42(id)|
# +-------------+
# |         true|
# +-------------+

否则直接调用注册副作用:

代码语言:javascript
复制
def my_udf(other_par, spark):
    def _(s):
        return len(s) == other_par
    name = "my_udf_{}".format(other_par)
    spark.udf.register(name, _, BooleanType())
    return udf(_, BooleanType())

my_udf_0 = my_udf(0, spark)

spark.sql("SELECT my_udf_0(array())").show()
# +-----------------+
# |my_udf_0(array())|
# +-----------------+
# |             true|
# +-----------------+

当然,像这样的简单操作不应该用udf完成,但我认为它只是一个玩具示例。如果不是,

代码语言:javascript
复制
from pyspark.sql.functions import size, length

size("some_col") == 42
length("some_col") == 42

是更好的选择。

票数 2
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/49821875

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档