PySpark UDF执行报错:TypeError: code()参数13需为str而非int
PySpark UDF执行时Python Worker抛出TypeError异常
错误详情
运行自定义PySpark UDF时,Python Worker抛出如下异常:
File "C:\PATH\SparkInstallation\spark-3.3.1-bin-hadoop3\python\lib\pyspark.zip\pyspark\serializers.py", line 471, in loads return cloudpickle.loads(obj, encoding=encoding) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ TypeError: code() argument 13 must be str, not int
用户排查后怀疑错误并非自身业务代码导致。
业务代码及预期
以下代码用于从推文数据集的tweet列中提取哈希标签,生成hashtags数组列:
from pyspark.sql.types import StringType, ArrayType import re from pyspark.sql import SparkSession import pyspark.sql.functions as pysparkfunctions spark = SparkSession.builder.appName('test').getOrCreate() savedTweets = spark.read.csv("testData/") def getHashtags(string): return re.findall(r"#(\w+)", string) getHashtagsUDF = pysparkfunctions.udf(getHashtags, ArrayType(StringType())) savedTweets = savedTweets.withColumn("hashtags", getHashtagsUDF(savedTweets['tweet'])) savedTweets.show()
UDF预期行为:
- 输入
" #a #b #c",输出['a', 'b', 'c'] - 输入
" a @b #c",输出['c']
错误原因及解决方案
这个错误本质是Python版本与Spark版本不兼容导致的序列化问题:
Spark 3.3.1依赖的cloudpickle版本,无法正确处理Python 3.10+版本中函数对象的code参数格式(Python 3.10调整了该参数的类型要求,从int改为str),从而触发类型不匹配异常。
可行解决方案:
- 降级Python版本:将Python版本降至3.9.x(Spark 3.3.1官方兼容Python 3.7/3.8/3.9)
- 升级Spark版本:将PySpark升级到3.4及以上版本,此类版本已修复对Python 3.10+的序列化兼容问题
- 调整UDF实现:避免在UDF中直接引用模块级函数(如
re.findall),可将逻辑内联,确保函数序列化时能正确捕获依赖
内容的提问来源于stack exchange,提问作者René Steeman
相关产品推荐
相关产品推荐

