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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:07:20