如何将自定义函数应用于PySpark DataFrame多列以生成新列
解决方案
完全不需要拼接单列,PySpark原生支持向UDF传入多列参数,函数内部可以直接获取每列的独立值,实现和Pandas apply类似的效果。
基础实现(兼容所有PySpark版本)
- 首先导入依赖
from pyspark.sql.functions import udf from pyspark.sql.types import BooleanType
- 定义自定义函数,参数顺序和后续传入的列顺序对应即可
def some_func(eventid, lat, lon): # 此处直接使用三个参数做运算,不需要自行拆分字段 # 示例逻辑可替换为实际运算规则 if lat is None or lon is None or eventid is None: return False return 30 < lat < 60 and 70 < lon < 140
- 注册UDF并应用到DataFrame
# 注册UDF,指定返回值类型为布尔型 verify_udf = udf(some_func, BooleanType()) # 直接按顺序传入需要的三个列,生成新列verified df = df.withColumn("verified", verify_udf("eventid", "lat", "lon"))
简化写法(PySpark 3.0及以上版本)
可以用装饰器直接注册UDF,减少冗余代码:
from pyspark.sql.functions import udf @udf(returnType=BooleanType()) def some_func(eventid, lat, lon): # 运算逻辑 if lat is None or lon is None or eventid is None: return False return 30 < lat < 60 and 70 < lon < 140 # 直接调用装饰后的函数即可 df = df.withColumn("verified", some_func("eventid", "lat", "lon"))
性能优化提示
如果数据量较大,推荐使用矢量化Pandas UDF,可以大幅降低Python UDF的序列化开销,处理效率更高:
import pandas as pd from pyspark.sql.functions import pandas_udf @pandas_udf(BooleanType()) def some_func(eventid: pd.Series, lat: pd.Series, lon: pd.Series) -> pd.Series: # 此处直接对pandas Series做矢量化运算,速度远快于逐行处理的普通UDF return (lat.between(30,60)) & (lon.between(70,140)) & (eventid.notna()) df = df.withColumn("verified", some_func("eventid", "lat", "lon"))
内容的提问来源于stack exchange,提问作者Siab Shafique
相关产品推荐
相关产品推荐

