PySpark中如何为UDF指定Double类型返回值?
解决PySpark UDF指定Double返回类型的问题
你的代码存在两个关键问题,修正后就能让final_value列成为Double类型:
- UDF参数不匹配:你定义的
lambda z: new_column(z)只传递了一个参数,但new_column需要接收row和year两个参数,会导致调用时报错。 - 返回类型写法错误:指定Double类型需要使用
DoubleType()(来自pyspark.sql.types模块),而非Double()。
修正后的完整代码
from pyspark.sql import functions as F from pyspark.sql.types import DoubleType # 该函数判断年份并对value_column进行乘法运算得到最终列 def new_column(row, year): if year == "2020": return row * 0.856 elif year == "2019": return row * 0.8566 else: return row # 直接关联函数并指定返回类型为DoubleType final_udf = F.udf(new_column, DoubleType()) # 或者用lambda传递两个参数的写法: # final_udf = F.udf(lambda val, yr: new_column(val, yr), DoubleType()) res = res.withColumn("final_value", final_udf(F.col('value_column'), F.col('year')))
关键说明
- 必须从
pyspark.sql.types导入DoubleType,这是PySpark用来标识数据类型的标准类,类似的还有StringType()、IntegerType()等。 - UDF定义时,要确保传入的参数数量和
new_column的参数数量一致,避免参数不匹配的错误。
内容的提问来源于stack exchange,提问作者ar_mm18
相关产品推荐
相关产品推荐

