PySpark根据边界列表为DataFrame列添加上下边界列异常排查
问题原因
你遇到的两个异常分别由以下原因导致:
- DataFrame无变化:PySpark所有DataFrame转换都是不可变操作,
withColumn不会修改原DataFrame,只会返回携带新列的新DataFrame,你没有将返回结果赋值给变量,原test对象自然不会有变化。 - UDF返回大量null:你声明UDF返回类型为
FloatType,但lambda逻辑中,当匹配到的边界是整数(比如0、7、41、235)时,Python返回的是int类型值,和声明的Float类型不匹配,Spark无法完成隐式转换就会返回null。 - 额外小问题:你的纯Python测试代码里存在笔误,
upper(I)中的I为大写,运行时会抛出NameError,改成小写i即可正常运行。
解决方法
方法1:修复原有UDF逻辑(改动最小)
只需要做两处修改:一是将UDF计算结果强制转为float匹配返回类型,二是将withColumn返回的新DataFrame赋值给变量:
from pyspark.sql import functions as F from pyspark.sql.types import FloatType border_list = [0, 1.5, 7, 41, 235] # 计算结果强制转float,匹配声明的返回类型 lower_border = F.udf(lambda x: float(max([i for i in border_list if x >= i])), FloatType()) upper_border = F.udf(lambda x: float(min([i for i in border_list if x < i])), FloatType()) # 必须将转换结果赋值给新变量 result = test.withColumn("LowBorder", lower_border("Population")) \ .withColumn("UpBorder", upper_border("Population")) result.show()
运行后输出完全符合预期:
+----------+---------+--------+ |Population|LowBorder|UpBorder| +----------+---------+--------+ | 1| 0.0| 1.5| | 2| 1.5| 7.0| | 234| 41.0| 235.0| | 0| 0.0| 1.5| | 6| 1.5| 7.0| | 7| 7.0| 41.0| | 35| 7.0| 41.0| | 46| 41.0| 235.0| | 8| 7.0| 41.0| +----------+---------+--------+
方法2:使用Spark原生函数(性能更优)
Python UDF需要在Python进程和JVM之间做数据序列化,性能远低于Spark内置函数,数据量大时推荐用原生算子实现,完全避免UDF开销:
from pyspark.sql import functions as F border_list = [0, 1.5, 7, 41, 235] border_col = F.array(*[F.lit(b) for b in border_list]) result = test.withColumn("border", border_col) \ .withColumn("LowBorder", F.array_max(F.expr("filter(border, x -> Population >= x)"))) \ .withColumn("UpBorder", F.array_min(F.expr("filter(border, x -> Population < x)"))) \ .drop("border") result.show()
该方法全程使用Spark内置算子,无跨进程序列化开销,执行效率更高,返回结果和方法1完全一致。
如果是固定分桶场景,也可以直接使用
pyspark.ml.feature.Bucketizer完成分桶后再关联边界值,边界数量较多时代码更简洁。
内容的提问来源于stack exchange,提问作者Bohdan Pylypenko
相关产品推荐
相关产品推荐

