Databricks中PySpark UDF声明式赋值触发ModuleNotFoundError求助
PySpark UDF声明式赋值报错ModuleNotFoundError: No module named 'src'(Databricks环境)
问题背景
在Databricks环境中使用PySpark时,采用非装饰器的声明式方式定义UDF会触发ModuleNotFoundError,但装饰器方式完全正常。报错核心是Python Worker节点无法找到src模块。
代码示例
正常工作的装饰器定义(src/lib/udfs.py)
# src/lib/udfs.py # Option 1 (works) from pyspark.sql.functions import udf from pyspark.sql.types import StringType @udf(returnType=StringType()) def apply_transformation(value: str) -> str: return value
报错的声明式赋值定义(src/lib/udfs.py)
# src/lib/udfs.py # Option 2 (doesn't work) from pyspark.sql.functions import udf from pyspark.sql.types import StringType def apply_transformation(value: str) -> str: return value udf_apply_transformation = udf(apply_transformation, StringType())
测试代码(tests/lib/test_udfs.py)
# 正常执行的测试 def test_udf(self): # 初始化DataFrame等逻辑 df.select(apply_transformation("value")).collect() # 报错的测试 def test_udf(self): # 初始化DataFrame等逻辑 df.select(udf_apply_transformation("value")).collect()
核心报错信息
pyspark.errors.exceptions.captured.PythonException: An exception was thrown from the Python worker. Please see the stack trace below. 'pyspark.serializers.SerializationError: Caused by Traceback (most recent call last): File "/databricks/spark/python/pyspark/serializers.py", line 189, in _read_with_length return self.loads(obj) File "/databricks/spark/python/pyspark/serializers.py", line 540, in loads return cloudpickle.loads(obj, encoding=encoding) ModuleNotFoundError: No module named 'src''.
问题原因
装饰器方式定义UDF时,PySpark内部会自动处理函数的序列化元数据,调整模块引用路径;而声明式赋值方式直接将原始函数传入udf(),cloudpickle序列化时会保留完整的src.lib.udfs.apply_transformation模块路径。但Databricks的Python Worker节点默认不会将项目根目录添加到Python路径,导致反序列化时无法解析src模块。
解决方案
方案1:将项目根目录添加到Python路径
在代码初始化阶段,手动将src所在的根目录加入Python路径:
import sys from pathlib import Path # 假设当前脚本位于src/lib,向上三级定位到项目根目录 project_root = Path(__file__).resolve().parent.parent.parent sys.path.append(str(project_root))
也可以通过Databricks集群配置全局设置,在集群的Spark配置中添加:
spark.driver.extraPythonPath /dbfs/your/project/root/path spark.executor.extraPythonPath /dbfs/your/project/root/path
方案2:修改函数的模块属性
手动重置函数的__module__属性,避免序列化时携带src路径:
# src/lib/udfs.py def apply_transformation(value: str) -> str: return value # 让序列化时使用全局模块标识,避免src路径被包含 apply_transformation.__module__ = "__main__" udf_apply_transformation = udf(apply_transformation, StringType())
方案3:调整Databricks作业目录配置
如果是通过Databricks作业运行代码,确保将src目录上传到DBFS,并在作业的任务设置中,将项目根目录设为工作目录,或者在任务脚本开头添加路径配置代码。
验证
应用任一方案后,重新运行测试代码,udf_apply_transformation即可正常执行,不再触发模块找不到的错误。
内容的提问来源于stack exchange,提问作者AlexWeh
相关产品推荐
相关产品推荐

