PySpark使用UDF映射批量添加列时返回空值的问题排查
PySpark使用UDF映射批量添加列时返回空值的问题排查
嗨,我一眼就发现问题出在哪了!你的UDF返回类型定义和实际返回值完全不匹配,这才导致所有新列都变成NULL啦。
问题根源分析
你写的test_map函数返回的是字符串'low'或'high',但注册UDF的时候却指定了返回类型是IntegerType()——Spark尝试把字符串转换成整数,自然转失败,结果就全是NULL了。
方案1:修正UDF的返回类型
只需要把UDF的返回类型改成StringType()就行,其他逻辑完全不变:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def test_map(col): if col < 5: score = 'low' else: score = 'high' return score # 关键修正:把IntegerType()改成StringType() test_udf = F.udf(test_map, StringType()) # 批量生成列映射的逻辑保持不变 mapp = {} for x in (1, 2): print(f'Now working {x}') mapp[f'limit_{x}'] = test_udf(F.col(f'int_{x}')) # 执行添加列操作 df.withColumns(mapp).show()
运行后就能得到你预期的结果:
+-----+-----+-------+-------+ |int_1|int_2|limit_1|limit_2| +-----+-----+-------+-------+ | 1| 3| low| low| | 2| 3| low| low| | 3| 5| low| low| | 4| 10| low| high| | 5| 20| high| high| +-----+-----+-------+-------+
方案2:放弃UDF,用Spark内置函数(更推荐)
其实这种简单的条件判断完全不需要写UDF,Spark内置的when/otherwise性能更好(UDF是Python黑盒,Spark无法对其进行优化,而内置函数能被Catalyst优化器处理),代码也更简洁:
from pyspark.sql import functions as F # 用内置函数定义映射逻辑 def get_limit_col(col_name): return F.when(F.col(col_name) < 5, 'low').otherwise('high') # 批量生成列映射 mapp = {} for x in (1, 2): mapp[f'limit_{x}'] = get_limit_col(f'int_{x}') df.withColumns(mapp).show()
这个版本不仅能得到同样的正确结果,在处理100列的大数据集时,性能会比UDF版本好很多。
另外你提到withColumns比循环withColumn快,这个思路完全正确——withColumns可以一次性提交所有列的逻辑,避免多次触发计划解析,确实更高效。
内容来源于stack exchange
相关产品推荐
相关产品推荐

