Databricks共享集群PySpark UDF无法导入Repos自定义模块问题
问题根因
单元格直接导入模块不报错,是因为Driver节点默认把当前Notebook所属Repo的根目录加进了Python搜索路径,但PySpark UDF是分发到Executor节点运行的,Executor的Python环境默认不会自动同步其他Repo的路径配置,反序列化UDF依赖时就会抛出ModuleNotFoundError: No module named 'test'错误。
解决方案(无需修改集群配置、无需打包egg/wheel)
全程只需要在Notebook最开头加几行配置代码,配置仅对当前Notebook会话生效,不会影响共享集群上其他用户的任务,完全适配快速迭代开发场景:
- 先拿到存放自定义模块的Repo在Databricks环境里的绝对路径
Databricks Repos的统一挂载根路径是/Workspace/Repos/,路径规则为/Workspace/Repos/<你的登录账号>/<存放模块的Repo名称>,比如你的登录账号是zhangsan@company.com,存test模块的Repo叫data-common,对应路径就是/Workspace/Repos/zhangsan@company.com/data-common。如果
test模块是Repo下的子文件夹,直接把路径定位到test文件夹的上一级目录即可。 - 在所有导入语句、业务代码之前插入如下代码,把路径同步给Driver和所有Executor:
import sys import os from pyspark.sql import SparkSession # 替换成你自己的目标路径 MODULE_PARENT_PATH = "/Workspace/Repos/替换成你的账号/替换成对应Repo名" # 给当前Driver节点添加搜索路径,保证本地导入正常 if MODULE_PARENT_PATH not in sys.path: sys.path.append(os.path.abspath(MODULE_PARENT_PATH)) # 把路径同步给所有Executor节点,不需要修改集群全局配置 spark = SparkSession.builder.getOrCreate() spark.sparkContext.addPyFile(MODULE_PARENT_PATH)
- 后续正常导入模块、定义UDF、调用
withColumn即可,不需要修改原有业务逻辑:
from test import pyspark_utils from pyspark.sql.functions import udf @udf("string") def process_udf(input_val): return pyspark_utils.do_process(input_val) # 此处调用不会再报模块找不到的错误 df = df.withColumn("processed_col", process_udf("raw_col"))
开发效率优化
如果需要频繁修改自定义模块代码,不需要重启集群,只要在Notebook开头加两行自动重载配置,修改完模块代码直接重新运行单元格即可生效:
%load_ext autoreload %autoreload 2
注意:直接传目录的方式在DBR 10.0及以上版本原生支持,如果你的集群用的是更早的版本,只需要把存放模块的文件夹压缩成zip包放在Repo里,把
addPyFile的参数换成zip包的绝对路径即可,不需要打标准egg/wheel包,也不需要做集群层面的上传安装操作。
内容的提问来源于stack exchange,提问作者bbl007
相关产品推荐
相关产品推荐

