使用PySpark --files选项在YARN集群模式下遇文件不存在错误求助
YARN集群模式下PySpark使用--files选项报「No such file or directory」错误排查
问题背景
在YARN集群模式下提交PySpark任务时,使用--files选项分发文件,运行时出现文件找不到的错误,无法访问指定文件。
示例输入文件(your_file_name.txt)
abc xyz pqr
代码文件(script.py)
from pyspark import SparkContext from pyspark.sql import SparkSession from pyspark import SparkFiles # 创建SparkSession spark = SparkSession.builder.getOrCreate() # 使用SparkFiles.get()访问文件 file_path = SparkFiles.get("your_file_name.txt") # 读取文件内容 with open(file_path, "r") as f: content = f.read() # 注意:此处存在缩进错误 # 处理文件内容 print(content)
提交命令
spark3-submit --master yarn --deploy-mode cluster --files your_file_name.txt script.py
报错日志
============================================================================================= LogType:stdout LogLastModifiedTime:Wed Dec 04 10:47:07 +0000 2024 LogLength:355 LogContents: Traceback (most recent call last): File "script.py", line 12, in with open(file_path, "r") as f: FileNotFoundError: [Errno 2] No such file or directory: '/dfs/10/yarn/nm/usercache/dmahajan/appcache/application_1724863745359_15570/spark-86bec527-b8d8-4ece-b449-36da0910eca7/userFiles-2056673d-bcef-4300-803a-62a368d6a82c/your_file_name.txt' End of LogType:stdout
排查与解决方案
1. 修复代码缩进错误
原代码中with open块内的content = f.read()未缩进,导致代码逻辑错误。修正后代码如下:
from pyspark import SparkContext from pyspark.sql import SparkSession from pyspark import SparkFiles # 创建SparkSession spark = SparkSession.builder.getOrCreate() # 使用SparkFiles.get()访问文件 file_path = SparkFiles.get("your_file_name.txt") # 读取文件内容 with open(file_path, "r") as f: content = f.read() # 修复缩进 # 处理文件内容 print(content)
2. 确认文件路径正确性
提交命令中--files your_file_name.txt需确保文件存在于提交节点的当前工作目录,若文件不在当前目录,需指定绝对路径(如--files /home/user/your_file_name.txt)。
3. 验证文件是否成功分发
通过YARN Web UI(通常为http://<resource-manager-host>:8088)找到对应应用,进入「Application Master」页面,查看「Files」标签,确认your_file_name.txt存在于userFiles目录下。若文件未上传,需检查提交命令是否正确,或文件是否有权限被读取。
4. 改用Spark分布式API读取文件
集群模式下,本地open()方法依赖节点本地文件,不符合Spark分布式架构设计。推荐使用Spark原生API读取文件,避免本地路径问题:
from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder.getOrCreate() # 用Spark API读取分布式文件 df = spark.read.text("your_file_name.txt") # 收集内容并格式化 content = "\n".join([row.value for row in df.collect()]) print(content)
内容的提问来源于stack exchange,提问作者Deepak Mahajan
相关产品推荐
相关产品推荐

