PySpark UDF使用ElementTree生成XML时报PicklingError如何解决
错误原因
- 你将
lxml.etree._Element类型的root变量定义在了UDF函数外部:Spark会对UDF进行序列化后分发到Executor节点执行,序列化时会捕获UDF闭包引用的所有变量,而lxml的元素对象本身不支持Python pickle序列化,直接触发报错。 - 原有代码还存在两处隐藏逻辑错误:
etree.SubElement第二个参数要求是字符串格式的标签名,你直接传入了整数类型的field1、field2会触发参数类型错误- 给节点
text属性赋值时要求传入字符串,直接传入整数类型的字段值也会报错
- 若
root定义在全局,还会出现多轮调用UDF时XML内容累加的问题,不符合每行生成独立XML的需求。
修复方案
把root元素的创建逻辑移到UDF内部,同时修正标签名、文本赋值的类型问题,并且显式指定UDF的返回类型为字符串类型即可正常运行,修改后代码如下:
import pyspark.sql.functions as F import pyspark.sql.types as T from lxml import etree data=[ (123,1,"string123","string456","string789")] importSchema=(T.StructType([ T.StructField("field1",T.IntegerType(),True), T.StructField("field2",T.IntegerType(),True), T.StructField("field3",T.StringType(), True), T.StructField("field4",T.StringType(),True), T.StructField("field5",T.StringType(),True) ])) df=spark.createDataFrame(data=data,schema=importSchema) def create_str(field1,field2,field3,field4,field5): # root元素创建移到UDF内部,不需要序列化传输 root = etree.Element('root') outer = etree.SubElement(root, 'outer') # 标签名使用字符串定义,如果需要用字段值做标签需转成字符串 field1s = etree.SubElement(outer, "field1") field2s = etree.SubElement(outer, "field2") field3s = etree.SubElement(outer, "field3") field4s = etree.SubElement(outer, "field4") field5s = etree.SubElement(outer, "field5") # 整数字段转成字符串后再赋值给text属性 field1s.text = str(field1) field2s.text = str(field2) field3s.text = field3 field4s.text = field4 field5s.text = field5 var=etree.tostring(root, pretty_print=True).decode('utf-8') return var # 显式指定UDF返回类型为字符串 udf_create_str = F.udf(create_str, T.StringType()) # 测试执行 df.withColumn("output", udf_create_str(df.field1,df.field2,df.field3,df.field4,df.field5)).show(truncate=False)
补充说明
运行后output列会生成标准XML字符串,不会再触发序列化错误,每行对应独立的XML结构,不会出现内容累加问题。如果需要用字段值作为标签名,直接把"field1"这类固定标签替换为str(field1)即可。
内容的提问来源于stack exchange,提问作者dcrowley01
相关产品推荐
相关产品推荐

