PySpark3.1.2内置函数注册及在spark.sql查询中使用的正确方法
PySpark 3.1.2 内置函数在spark.sql中使用的正确方案
报错原因
你当前的报错根源是:pyspark.sql.functions下的内置函数本质是Column表达式处理函数,输入输出都要求是Column类型,不能直接作为UDF注册给SQL层调用——SQL层传递给UDF的是行内的具体字面量值,不是Column对象,所以传入整数-5的时候会触发类型不匹配报错。
报错原文:
...TypeError: Invalid argument, not a string or column: -5 of type <class 'int'>. For column literals, use 'lit', 'array', 'struct' or 'create_map' function...
正确实现方式
方式1:直接使用Spark SQL原生内置函数(最优)
Spark SQL本身已经内置了几乎所有pyspark.sql.functions下的同名函数,不需要额外注册即可直接使用,你的纯SQL写法本来就是正确的:
# 无需注册任何函数,直接调用SQL原生ABS即可 rv1 = spark.sql('SELECT ABS(x + y) AS z FROM df').toPandas()
如果需要给内置函数自定义别名,可以直接创建SQL函数别名:
spark.sql('CREATE FUNCTION abs_builtin AS org.apache.spark.sql.catalyst.expressions.Abs')
之后就可以直接在SQL里用abs_builtin别名调用绝对值函数了。
方式2:封装为Python UDF后注册(适用于需要自定义改造逻辑的场景)
如果你确实需要把Python侧的逻辑封装成UDF给SQL调用,要写接收普通值的处理函数,不能直接传F.abs:
# 直接封装Python原生abs函数注册,也可以替换为你自己的自定义逻辑 spark.udf.register('abs_builtin', lambda x: abs(x) if x is not None else None, T.LongType()) rv2 = spark.sql('SELECT abs_builtin(x + y) AS z FROM df').toPandas()
性能说明
优先使用方式1调用原生SQL内置函数,性能比自定义UDF高10~100倍,没有Python-JVM跨进程通信的额外开销。
内容的提问来源于stack exchange,提问作者Russell Burdt
相关产品推荐
相关产品推荐

