Databricks中UDF触发PicklingError:无法序列化非读模式文件
问题重现
我在Databricks上遇到如下错误:
/databricks/spark/python/pyspark/cloudpickle/cloudpickle_fast.py in dumps(obj, protocol, buffer_callback) 71 file, protocol=protocol, buffer_callback=buffer_callback 72 ) ---> 73 cp.dump(obj) 74 return file.getvalue() 75 /databricks/spark/python/pyspark/cloudpickle/cloudpickle_fast.py in dump(self, obj) 600 def dump(self, obj): 601 try: --> 602 return Pickler.dump(self, obj) 603 except RuntimeError as e: 604 if "recursion" in e.args[0]: /databricks/spark/python/pyspark/cloudpickle/cloudpickle_fast.py in _file_reduce(obj) 314 ) 315 if "r" not in obj.mode and "+" not in obj.mode: --> 316 raise pickle.PicklingError( 317 "Cannot pickle files that are not opened for reading: %s" 318 % obj.mode PicklingError: Cannot pickle files that are not opened for reading: a
背景:我尝试把打印语句的输出记录到文件里,但代码进入UDF函数时触发了这个错误。该操作在部分Pandas UDF中可以正常运行,但在另一些中却不行。
问题根源
Spark执行UDF时,会把UDF相关的所有对象通过cloudpickle序列化,然后分发到集群的各个工作节点执行。但cloudpickle明确禁止序列化非只读模式的文件句柄(这里你的文件是"a"追加写入模式):
- 文件句柄和本地进程强绑定,就算序列化到其他节点也无法正常使用;
- 写入模式的文件序列化存在数据竞争、文件损坏的风险,因此被直接拦截。
至于部分UDF能正常运行,是因为那些UDF的文件操作没有把文件对象纳入UDF的序列化上下文——比如在UDF内部临时打开后立刻关闭,或者文件对象没有被UDF的闭包、参数引用到,序列化时不会处理它。而报错的UDF里,你在UDF外部打开了文件,又在UDF内部引用了这个文件对象,导致序列化时尝试处理这个写入模式的句柄。
解决方案
不要在UDF外部持有写入模式的文件句柄
把文件打开、写入、关闭的逻辑放到UDF内部,但要注意:集群每个节点的本地文件系统是独立的,这样每个节点会写到自己本地的文件,不是同一个统一文件。如果需要统一存储,改用DBFS路径,并且通过Spark的分布式API写入(比如将数据收集到Driver后统一写入,或者用DataFrame的写入接口)。示例代码:
from pyspark.sql.functions import pandas_udf import pandas as pd @pandas_udf("string") def my_udf(data: pd.Series) -> pd.Series: # 在UDF内部完成文件操作,写完立刻关闭 with open("/dbfs/path/to/your/log.txt", "a") as f: f.write(f"Processing batch: {data.head()}\n") return data.str.upper()改用日志框架替代自定义文件写入
不要自己写文件捕获print输出,直接配置Spark的log4j日志,把UDF的输出定向到指定的日志文件,或者利用Databricks自带的日志收集功能(比如查看集群节点日志、作业日志),这才是分布式环境的最佳实践。避免序列化文件对象
如果必须使用文件,只把文件路径作为参数传入UDF,在UDF内部完成文件操作,确保文件句柄不会被UDF的上下文捕获,也就不会被序列化。
内容的提问来源于stack exchange,提问作者Pranav Gupta

