You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 20:52:06