Spark调用pipe时SparkFiles.get()返回Driver而非Worker路径问题
根本原因
这个问题本质是对Spark的执行边界和API逻辑理解存在偏差,核心是两个误区:
SparkFiles.get()没有全局统一返回值:这个方法的返回值完全取决于调用它的进程位置,在Driver进程里调用,返回的就是Driver节点本地的Spark文件缓存绝对路径;在Worker节点的Executor进程里调用,才会返回对应Worker本地的缓存路径。你在Driver端拼接命令时调用该方法,拿到的天生就是Driver专属路径,这个字符串序列化发到Worker后,Worker本地根本没有同名的目录结构,必然报文件不存在。RDD.pipe()不会做路径自动替换:你传给pipe()的命令字符串会被原样序列化下发到Executor执行,框架不会自动识别字符串里的Driver侧路径,也不会帮你替换成Worker本地的对应路径。
容易漏判的触发点
你用的是DBFS上的源文件,sparkContext.addFile()把文件分发到Worker节点后,默认不会给文件赋予可执行权限。在很多Linux环境下,调用无执行权限的文件时,抛出的异常是FileNotFound而非Permission denied,非常容易误导排查方向。
另外你贴的复现代码里还有个低级笔误:定义的RDD变量名叫pipe_rdd,最后collect时调用的是pipe_tokenised_rdd,这个变量根本没定义,改完路径问题也要注意修正这个笔误。
修复方案
不要在Driver端提前拼接带绝对路径的命令,选下面任意一种方式改即可:
- 直接用文件名调用可执行文件
Spark启动pipe子进程时,会自动把addFile分发文件的缓存目录加到子进程的PATH环境变量里,不需要写全路径,直接用文件名就能找到:spark.sparkContext.addFile("/dbfs/FileStore/Custom_Executable") files_rdd = spark.sparkContext.parallelize(files_list) # 直接用文件名,不要调用Driver端的SparkFiles.get()拼绝对路径 cmd = "Custom_Executable the-function-name --to company1 --from company2 -input - -output -" pipe_rdd = files_rdd.pipe(cmd, env={'SOME_ENV_VAR': env_var_val}) print(pipe_rdd.collect()) - 需要配置文件权限/必须用绝对路径时,在Executor端完成初始化
如果你需要手动给可执行文件加权限,或者有特殊需求必须指定绝对路径,通过mapPartitions在Executor进程内完成路径获取和权限配置即可:from pyspark import SparkFiles import os spark.sparkContext.addFile("/dbfs/FileStore/Custom_Executable") files_rdd = spark.sparkContext.parallelize(files_list) # 每个分区在Executor上启动时先初始化环境 def init_partition(iterator): exe_local_path = SparkFiles.get("Custom_Executable") # 给可执行文件加可执行权限 os.chmod(exe_local_path, 0o755) return iterator files_rdd = files_rdd.mapPartitions(init_partition) # 权限配置完成后直接用文件名调用即可 pipe_rdd = files_rdd.pipe("Custom_Executable the-function-name --to company1 --from company2 -input - -output -", env={'SOME_ENV_VAR': env_var_val}) print(pipe_rdd.collect())
内容的提问来源于stack exchange,提问作者mikelus
相关产品推荐
相关产品推荐

