PySpark调用字符串拼接UDF报TypeError: decoding str is not supported
报错原因
这个类型错误是UDF调用时参数类型不匹配+用法不规范导致的,核心问题有两个:
- 传参类型不符合PySpark UDF要求:在DataFrame DSL语法中调用UDF时,所有传入参数必须是
Column类型——要么是通过col()引用的DataFrame列,要么是通过lit()包装的常量值。原代码里直接把"firstname"、"lastname"两个裸Python字符串作为参数传入,Spark在对参数做序列化、解码处理时无法识别这种非Column类型的字符串值,就会抛出decoding str is not supported的类型错误。 - UDF调用方式不规范:
spark.udf.register()的核心作用是把Python函数注册到Spark SQL的函数注册表,供写SQL查询时调用。如果要在select、withColumn这类DataFrame原生方法里直接传UDF调用逻辑,必须接住register()方法返回的UDF实例,否则当前代码作用域里根本不存在可调用的stringConcat_udf对象,也会触发执行异常。
修复方案
给两种可直接运行的实现方式:
方式1:DataFrame DSL 调用
注册UDF时接住返回的可调用对象,所有参数统一转为Column类型即可:
from pyspark.sql.functions import lit, col def stringConcat(separator: str, first: str, second: str): return first + separator + second # 接收注册后返回的UDF实例,供DSL语法直接调用 stringConcat_udf = spark.udf.register("stringConcat_udf", stringConcat) customerDf.select( "firstname", "lastname", # 列引用用col()包装为Column类型,分隔符常量用lit()包装 stringConcat_udf(lit("-"), col("firstname"), col("lastname")).alias("full_name") ).show()
方式2:SQL表达式调用
注册UDF后用expr()包裹SQL风格的调用语句,不需要手动包装Column类型,写法更简洁:
from pyspark.sql.functions import expr def stringConcat(separator: str, first: str, second: str): return first + separator + second spark.udf.register("stringConcat_udf", stringConcat) customerDf.select( "firstname", "lastname", # SQL表达式中常量用单引号包裹,列名直接书写即可 expr("stringConcat_udf('-', firstname, lastname) as full_name") ).show()
额外优化提示:这个分隔符拼接场景不需要自定义Python UDF,Spark内置的concat_ws函数原生支持该功能,性能比Python UDF高一个数量级(避免了Python和JVM之间的跨进程序列化开销),直接调用concat_ws(lit("-"), col("firstname"), col("lastname"))就能实现相同效果。
内容的提问来源于stack exchange,提问作者QueryQuasar
相关产品推荐
相关产品推荐

