Azure Databricks中PySpark自定义UDF报ModuleNotFoundError求助
问题:PySpark UDF执行时出现ModuleNotFoundError但导入正常
已参考相关问题但未解决,以下是具体场景:
代码仓库结构
|-run_pipeline.py |-__init__.py |-data_science |--__init__.py # 注意:原结构中写的__init.py__是错误的,需修正为标准的__init__.py |--text_cleaning |---text_cleaning.py |---__init__.py # 同上,需修正文件名
代码实现
run_pipeline notebook 代码
from data_science.text_cleaning import text_cleaning path = os.path.join(os.path.dirname(__file__), os.pardir) sys.path.append(path) spark = SparkSession.builder.master( "local[*]").appName('workflow').getOrCreate() df = text_cleaning.basic_clean(spark_df)
text_cleaning.py 代码
def basic_clean(df): print('Removing links') udf_remove_links = udf(_remove_links, StringType()) df = df.withColumn("cleaned_message", udf_remove_links("cleaned_message")) return df
报错信息
执行df.show()时抛出错误:
Exception has occurred: PythonException (note: full exception trace is shown but execution is paused at: <module>) An exception was thrown from a UDF: 'pyspark.serializers.SerializationError: Caused by Traceback (most recent call last): File "/databricks/spark/python/pyspark/serializers.py", line 165, in _read_with_length return self.loads(obj) File "/databricks/spark/python/pyspark/serializers.py", line 466, in loads return pickle.loads(obj, encoding=encoding) ModuleNotFoundError: No module named 'data_science''. Full traceback below: Traceback (most recent call last): File "/databricks/spark/python/pyspark/serializers.py", line 165, in _read_with_length return self.loads(obj) File "/databricks/spark/python/pyspark/serializers.py", line 466, in loads return pickle.loads(obj, encoding=encoding) ModuleNotFoundError: No module named 'data_science'
核心疑问
为何模块导入能正常工作,但执行UDF时却出现找不到模块的错误?
问题原因与解决办法
原因
PySpark采用Driver-Worker分布式执行架构:
- 模块导入在Driver节点完成,你添加的
sys.path仅对Driver生效,但Worker节点的Python环境无法识别该路径 - UDF代码会被序列化后分发到Worker节点执行,Worker在反序列化时找不到
data_science模块,因为其Python路径未包含你的代码目录 - 代码中先导入模块再添加路径,虽然Driver可能因当前目录或缓存侥幸识别到模块,但Worker完全无法访问
解决步骤
- 修正文件名错误:把仓库中所有
__init.py__重命名为__init__.py,这是Python识别包的必要条件 - 调整路径添加顺序:在导入模块前先把代码根目录添加到
sys.path,确保Driver能正确识别模块:
import sys import os from pyspark.sql import SparkSession # 先获取代码根目录并添加到路径 root_path = os.path.dirname(os.path.abspath(__file__)) sys.path.append(root_path) # 再导入模块 from data_science.text_cleaning import text_cleaning spark = SparkSession.builder.master("local[*]").appName('workflow').getOrCreate() df = text_cleaning.basic_clean(spark_df)
- 同步Worker节点的模块路径:在Databricks环境中,需确保所有Worker节点都能访问到你的代码:
- 方法一:安装为集群库:将代码打包成wheel文件,上传到Databricks并安装为集群级别的库,所有节点会自动识别模块
- 方法二:分发代码目录:用
dbutils将代码目录复制到DBFS公共路径,再添加到sys.path并通知Spark分发模块:# 复制本地代码到DBFS(假设本地代码在Workspace的/Repos/your-repo路径) dbutils.fs.cp("file:/Workspace/Repos/your-repo/", "dbfs:/path/to/your/repo", recurse=True) # 添加DBFS路径到sys.path sys.path.append("/dbfs/path/to/your/repo") # 分发模块文件到Worker节点 spark.sparkContext.addPyFile("/dbfs/path/to/your/repo/data_science/text_cleaning/text_cleaning.py")
内容的提问来源于stack exchange,提问作者user139442
相关产品推荐
相关产品推荐

