PySpark多参数UDF传字符串参数被误识别为列名报错如何解决
问题原因
PySpark调用自定义UDF时,默认会将传入的字符串参数解析为同名列的引用,而非静态字符串字面量。你代码中直接传递"name"作为UDF的第二个入参时,Spark会尝试在当前DataFrame的列中查找名为name的列,由于你的表列仅包含[zipcode, customer_id, customfields, timestamp, essential_item],不存在该列,就触发了如下列解析异常:
Exception: cannot resolve '`name`' given input columns: [zipcode, customer_id, customfields, timestamp, essential_item];;
解决方案
你需要使用PySpark内置的lit()函数将静态字面量封装为Column类型的字面量对象后再传入UDF,具体修改步骤如下:
- 导入
lit函数
from pyspark.sql.functions import lit
- 修改UDF调用代码,将字符串参数用
lit()包裹
修正后的完整可运行代码如下:
# rec 是 DataFrame from pyspark.sql.functions import from_json, lit, udf from pyspark.sql.types import StringType def findValue(records, name): for item in records: if item.key == name: return item.value return None findValueUDF = udf(lambda x,y: findValue(x,y), StringType()) # 核心修改:将原来的"name"替换为lit("name") rec = rec.withColumn("name", findValueUDF(from_json(rec.customfields, schema), lit("name")))
扩展说明:不止字符串,数值、布尔值等所有非列引用的静态参数传入UDF时,都建议用
lit()包裹,避免Spark解析逻辑导致的异常。
内容的提问来源于stack exchange,提问作者Miroslav Petrovic
相关产品推荐
相关产品推荐

