在Dataproc/Yarn运行Spark时临时文件报FileNotFoundException的咨询
Spark处理二进制文件的推荐方案及临时文件存储建议
一、临时文件存储的问题修复与替代路径
- Dataproc节点
/tmp的局限性:Dataproc集群中/tmp目录有自动清理机制,且Driver节点创建的临时目录无法被Worker节点跨访问(本地运行时Driver与Worker同节点,无此问题),同时可能存在权限限制,这是你遇到FileNotFoundException的核心原因。 - 推荐的临时存储路径:
- 使用节点专属的系统临时目录:通过
System.getProperty("java.io.tmpdir")动态获取每个Worker节点的本地临时目录,每个Executor会在自己节点的专属路径操作,避免跨节点访问问题。 - 改用
/var/tmp目录:该目录清理周期更长,权限配置更宽松,适合存放转换过程中的临时文件。
- 使用节点专属的系统临时目录:通过
- 强制清理临时文件:转换完成后必须显式删除临时文件/目录,避免占用节点存储,可使用
Files.deleteIfExists()或File.delete()方法处理。
二、分布式处理二进制文件的正确流程
- 在Executor节点本地执行全流程:
不要在Driver节点做下载、转换操作(单点无法利用集群并行能力,且Worker无法访问Driver本地文件),应把逻辑放到map/flatMap等分布式算子中,让每个Executor处理自己分区的文件。示例Java代码思路:// 从分布式存储读取待处理文件的路径列表 JavaRDD<String> filePaths = sparkContext.textFile("gs://your-bucket/file-list.txt"); filePaths.map(filePath -> { // 1. 下载二进制文件到Executor本地临时目录 String tempDir = System.getProperty("java.io.tmpdir"); Path localInputFile = Paths.get(tempDir, new File(filePath).getName()); // 实现从存储系统(如GCS/HDFS)下载文件到localInputFile的逻辑 // 2. 调用外部库执行格式转换,传入本地文件路径 String convertedFilePath = externalLibrary.convert(localInputFile.toString()); // 3. 将转换后的文件上传回分布式存储 uploadToStorage(convertedFilePath, "gs://your-bucket/converted/" + new File(filePath).getName()); // 4. 清理本地临时文件 Files.deleteIfExists(localInputFile); Files.deleteIfExists(Paths.get(convertedFilePath)); return filePath + " 转换完成"; }).collect(); - 使用节点本地磁盘优化大文件处理:如果临时文件体积较大,可使用Dataproc节点的本地磁盘(路径一般为
/mnt/local-disk),需提前通过集群配置或初始化脚本挂载并设置权限,能获得更好的IO性能,也避免临时目录的自动清理问题。
三、权限与访问配置
- 确保Executor进程权限:Dataproc的Executor默认以
yarn用户运行,需确保临时目录对yarn用户有读写权限,可通过初始化脚本配置:sudo mkdir -p /mnt/local-disk/temp sudo chown yarn:yarn /mnt/local-disk/temp sudo chmod 755 /mnt/local-disk/temp - 杜绝跨节点文件依赖:绝对不要让Worker节点访问Driver节点创建的本地文件,分布式场景下必须通过GCS、HDFS等分布式存储传递文件。
内容的提问来源于stack exchange,提问作者Allan Silva
相关产品推荐
相关产品推荐

