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

Databricks Connect环境下运行PySpark UDF出现ModuleNotFoundError的问题求助

解决Databricks Connect下PySpark UDF的ModuleNotFoundError问题

你的问题核心在于Databricks Connect的架构差异:本地运行时,Driver和Executor在同一个环境,自定义模块可以直接被访问;但Databricks Connect模式下,Driver运行在你的本地机器,而Executor是Databricks集群上的Worker节点——这些Worker并没有你的demo模块,所以当UDF被序列化发送到Worker执行时,就会抛出ModuleNotFoundError。

下面是几个可行的解决方案:

方案1:将自定义模块部署到Databricks集群

你需要把包含MyClass的demo模块同步到集群,让所有Worker节点都能访问到:

  • 方法A:用Databricks CLI上传模块到DBFS
    先把本地的demo目录上传到DBFS指定路径:
    databricks fs cp -r ./demo dbfs:/custom-modules/demo
    
    然后在初始化SparkSession时,给Worker节点的PYTHONPATH添加这个路径:
    spark = SparkSession.builder \
        .config("spark.executorEnv.PYTHONPATH", "/dbfs/custom-modules") \
        .getOrCreate()
    
  • 方法B:打包成Wheel包安装为集群库
    把demo模块打包成标准的Wheel包:
    python setup.py bdist_wheel
    
    然后在Databricks工作区中创建Python库,上传这个Wheel包并安装到你的集群,这样所有Worker都会自动加载这个模块。

方案2:重构代码,避免序列化整个类实例

如果不想部署模块,可以把类的逻辑拆解,避免传递整个MyClass实例到Worker:
比如把核心逻辑改成纯函数,用广播变量传递类实例中的固定参数(比如self.number):

import pyspark.sql.functions as f
from pyspark.sql import SparkSession, Row
from pyspark.sql import types as t

def calculate_total(base_num, num1, num2):
    return str(base_num + num1 + num2)  # 注意原代码遗漏了return,这里补上

spark = SparkSession.builder.getOrCreate()
# 广播固定的base_num,避免序列化整个类
broadcast_base = spark.sparkContext.broadcast(10)

# 定义不依赖本地模块的UDF
@f.udf('string')
def total_udf(num1, num2):
    return calculate_total(broadcast_base.value, num1, num2)

# 数据处理部分
rdd = spark.sparkContext.parallelize(
    [
        Row(num1=1, num2=1),
        Row(num1=1, num2=2),
        Row(num1=2, num2=2),
        Row(num1=2, num2=3),
    ]
)
schema = t.StructType(
    [
        t.StructField("num1", t.IntegerType(), True),
        t.StructField("num2", t.IntegerType(), True)
    ]
)
input_df = spark.createDataFrame(rdd, schema)
result = input_df.withColumn("new_col", total_udf("num1", "num2"))
result.show()

方案3:启用本地依赖同步(Databricks Runtime 10.4+)

如果你的集群使用Runtime 10.4或更高版本,可以开启Databricks Connect的本地依赖同步功能,它会自动把本地的自定义模块同步到集群Worker:

spark = SparkSession.builder \
    .config("spark.databricks.connect.localDependencySync.enabled", "true") \
    .getOrCreate()

这个方法适合小型模块,大模块还是建议用方案1更稳定。

为什么你的修改没生效?

你调整后的take_df方法仍然依赖MyClass实例,当UDF被序列化发送到Worker时,需要反序列化这个实例——但Worker没有demo模块,所以还是会报错。只有让Worker能访问到模块,或者彻底避免传递类实例,才能解决问题。

内容的提问来源于stack exchange,提问作者Kasia Kulma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 10:32:46