纯Spark SQL环境下能否使用Spark Pandas UDF及对应实现方法
问题结论
该调用方式完全可行,按照以下步骤操作即可:
操作步骤
- 首先将你写好的包含UDF定义和注册逻辑的Python代码保存为独立文件,比如命名为
udf_defs.py,文件内容如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import pandas_udf import pandas as pd spark = SparkSession.builder.getOrCreate() @pandas_udf(returnType="long") def add_one(v: pd.Series) -> pd.Series: return v.add(1) spark.udf.register("add_one", add_one)
- 调用
spark-sql命令时通过--py-files参数传入上述Python文件,让Spark会话启动时自动加载注册好的UDF,完整命令如下:
spark-sql --py-files udf_defs.py -e 'select add_one(1)'
注意事项
- 本地运行时需确保本地Python环境安装了匹配版本的
pandas、pyarrow依赖;集群运行时需保证所有Worker节点的Python环境都安装了上述依赖,否则pandas UDF会执行报错 - 如果需要无需每次传参即可调用该UDF,可以将UDF注册为Hive元数据的永久函数,注册时将Python文件上传到HDFS等集群可访问的统一路径即可
内容的提问来源于stack exchange,提问作者Neil McGuigan
相关产品推荐
相关产品推荐

