Spark(PySpark on YARN)中间文件的磁盘存储位置查询
核心结论
YARN模式下,Spark的临时文件(shuffle输出、persist磁盘文件)不会使用spark.local.dir或SPARK_LOCAL_DIRS配置,而是完全依赖YARN NodeManager分配的容器专属本地目录,且任务结束后会被自动清理。
具体存储路径定位
基础路径来源
查看yarn-site.xml中的yarn.nodemanager.local-dirs配置(比如你提到的hadoop/yarn/nm-local-dir),但实际文件存储在该路径的深层子目录中:<yarn.nodemanager.local-dirs>/hadoop/yarn/nm-local-dir/usercache/<你的用户名>/appcache/<应用ID>/container_<容器ID>/spark-local-*其中
spark-local-*是Spark为每个容器创建的临时存储文件夹。为什么之前没找到?
- YARN在Spark任务执行完成后,会自动清理对应容器的本地临时目录,所以必须在任务运行期间去工作节点上查找。
- 你测试代码中的
count()执行完毕后,若没有后续阻塞操作,Spark会自动释放persist的缓存文件,导致文件被删除。
实时定位方法
修改日志配置获取路径
在$SPARK_HOME/conf/log4j.properties中添加:log4j.logger.org.apache.spark.storage.BlockManager=INFO重启任务后,在Container日志中搜索
BlockManager关键词,会打印出类似如下的日志,直接显示本地存储目录:16/04/07 15:37:20 INFO BlockManager: Using local directory: /hadoop/yarn/nm-local-dir/usercache/xxx/appcache/application_12345/container_12345/spark-local-1234567890
通过YARN UI定位
打开YARN ResourceManager UI(默认端口8088),找到你的应用,进入对应Container的日志页面,直接搜索上述关键词即可获取路径。
测试代码优化(方便排查)
修改你的测试代码,添加阻塞逻辑让任务保持运行,以便有足够时间去节点上查找文件:
from pyspark import storagelevel df_sales = spark.read.load("gs://monsoon-credittech.appspot.com/spark_datasets/sales_parquet") df_products = spark.read.load("gs://monsoon-credittech.appspot.com/spark_datasets/products_parquet") df_merged = df_sales.join(df_products, df_sales.product_id == df_products.product_id, 'inner') df_merged.persist(storagelevel.StorageLevel.DISK_ONLY) df_merged.count() # 添加阻塞,保持任务运行 input("Press Enter to exit...")
内容的提问来源于stack exchange,提问作者figs_and_nuts

