YARN模式下SparkSession调用addFile无法将文件上传至Executor节点
问题修复方案:将HDFS文件同步到所有Spark Executor
你的问题核心是**SparkFiles.get()的调用位置错误**——它需要在Executor执行的分布式任务代码中使用,而非Driver本地代码里直接打印。
原因解释
spark.sparkContext.addFile()确实会将HDFS上的文件分发到所有Executor,但这个分发动作是延迟执行的,只有当第一个分布式任务(比如RDD的map/foreach操作)启动时,才会把文件复制到Executor的工作目录。- Driver端调用
SparkFiles.get()只会返回Driver本地缓存的文件路径,和Executor上的路径完全无关,所以你看到的路径自然只在Driver存在。
修改后的代码示例
from pyspark.sql import SparkSession from pyspark import SparkFiles import os if __name__ == "__main__": spark = SparkSession.builder.enableHiveSupport()\ .master("yarn").getOrCreate() # 先添加文件,确保在创建分布式任务前执行 spark.sparkContext.addFile("hdfs:///tmp/dummy.txt") # 通过分布式任务触发文件分发并检查Executor上的文件 def check_file_on_executor(): file_path = SparkFiles.get("dummy.txt") exists = os.path.exists(file_path) return f"Executor节点文件路径: {file_path}, 文件存在: {exists}" # 创建一个空RDD来触发分布式任务(实际场景中可替换为你的业务RDD/DataFrame) result = spark.sparkContext.parallelize(range(2)).map(lambda x: check_file_on_executor()).collect() # 打印各Executor返回的结果 for res in result: print(res) spark.stop()
关键注意事项
- 调用时机:
addFile必须在所有分布式任务(RDD/DataFrame操作)之前执行,否则文件不会被分发。 - 文件路径唯一性:每个Executor上的文件路径会不同,但文件名保持一致,通过
SparkFiles.get("dummy.txt")可以统一获取到当前Executor上的正确路径。 - 大文件场景:如果文件很大,不建议用
addFile(会给每个Executor复制一份,占用大量磁盘),可以直接让任务从HDFS读取文件路径hdfs:///tmp/dummy.txt。
内容的提问来源于stack exchange,提问作者Felix
相关产品推荐
相关产品推荐

