我知道我可以注册一个UDFand函数,因为它可以在SQL查询中使用:
def example(s):
return len(s)
sqlContext.udf.register("example_udf", example)
spark.sql("SELECT example_udf(col) FROM data")或者我可以用udf包装Python函数,这样就可以将它应用于dataframe:
from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType
example_udf = udf(example)
data.select(example_udf('col'))在我的例子中,由于我需要向UDF传递一些其他参数,所以我为UDF构建了一个嵌套函数:
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。
发布于 2018-04-13 17:08:31
如果您使用的是最新的和最大的(2.3),就不要直接使用udf:
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|
# +-------------+否则直接调用注册副作用:
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完成,但我认为它只是一个玩具示例。如果不是,
from pyspark.sql.functions import size, length
size("some_col") == 42
length("some_col") == 42是更好的选择。
https://stackoverflow.com/questions/49821875
复制相似问题