PySpark+Zeppelin:集群本地自定义库导入失效,sc.addFile()却生效?
这其实是PySpark的Driver-Executor架构和模块导入机制共同导致的问题,咱们一步步拆解:
核心原因:Driver与Executor的sys.path是完全独立的
你在第一种方法里执行的sys.path.insert(0, "/opt/repo/folder/"),只修改了Driver进程的sys.path——也就是Zeppelin所在节点运行的PySpark驱动进程的路径。而集群里真正执行RDD分区计算的Executor节点进程,它们的sys.path完全没被改动。
当你调用function(arguments)处理RDD时,PySpark会把function通过pickle序列化后发送到各个Executor节点执行。Executor在反序列化这个函数后,需要重新导入它依赖的module模块,但此时Executor的sys.path里并没有/opt/repo/folder/——哪怕你已经把仓库克隆到了所有节点的这个路径也没用,因为Executor进程启动时默认不会自动把这个路径加入sys.path。
这就是错误会在pickle.loads环节触发的原因:Executor找不到要导入的模块,直接抛出ImportError。
第二种方法为什么能正常工作?
sc.addFile("/opt/repo/folder/module.py")会告诉Spark主动把指定文件分发到所有Executor节点的临时工作目录(也就是SparkFiles.getRootDirectory()指向的路径)。之后你把这个临时目录加入sys.path,不管Executor节点原本有没有这个文件,此时它都能在本地找到module.py,自然就能成功导入模块并执行函数了。
额外补充:PySpark的序列化逻辑
PySpark在传递函数到Executor时,只会序列化函数本身,不会序列化它依赖的整个模块。它默认假设Executor端能找到这些依赖:要么是已经安装在所有节点的全局库,要么是通过addFile/addPyFile主动分发到Executor的文件。你的第一种方法虽然把文件放到了所有节点,但没有让Executor的sys.path包含这个路径,所以依然找不到模块。
内容的提问来源于stack exchange,提问作者kingledion

