Spark Scala中No space left on device错误的排查与解决
问题根源解析
Shuffle文件的管理逻辑
Spark的Shuffle临时文件由Executor的BlockManager负责管理,只有当关联的RDD/DataSet的所有依赖计算完成,且Spark确认这些数据不会被再次复用的情况下,才会触发自动清理。在你的循环迭代场景中,所有任务共享同一个SparkSession,Spark不会将每次迭代视为独立的上下文,因此第一次迭代生成的Shuffle文件会被BlockManager保留,不会自动清理——这不属于默认的自动清理场景。手动删除引发异常的原因
你直接删除/tmp/spark-local-dir-1下的文件后,Spark的BlockManager仍保留着这些文件的元数据记录。后续迭代任务执行时,BlockManager会尝试访问这些已被删除的块,导致文件未找到异常,同时因为元数据与实际磁盘文件不一致,进一步加剧了磁盘空间的管理混乱。
可行的解决方案
1. 每次迭代后显式清理临时数据
在每次迭代的处理逻辑完成后,显式将当前迭代生成的DataFrame/RDD标记为可清理,让BlockManager主动删除对应的Shuffle文件。注意不要清理你需要长期缓存的filterDF和noFilterDF:
fileList.sliding(500).foreach(x => { val data = SparkSession.read.format("avro").option("header", "true").load(x: _*) val joinedData = data.join(filterDF).join(noFilterDF) // 修正变量名拼写错误 val results = joinedData.groupBy("x","y").agg(sum("p")) results.write.mode(SaveMode.Overwrite).parquet(filePath) // 显式清理当前迭代生成的临时数据 data.unpersist(true) joinedData.unpersist(true) results.unpersist(true) })
同时修正代码中的拼写错误:filterDF.cahce()改为filterDF.cache(),nofilterDF统一为noFilterDF,ist(...)改为List(...)。
2. 调整Spark配置优化Shuffle文件管理
通过以下配置调整,让Spark更主动地清理Shuffle文件并优化磁盘使用:
- 启用外部Shuffle服务:
外部Shuffle服务会独立管理Shuffle文件,避免Executor退出时的文件残留,同时更智能地清理不再使用的Shuffle数据。spark.shuffle.service.enabled=true - 调整Shuffle文件清理阈值:
开启检查点清理并缩短GC跟踪的间隔,让ContextCleaner更频繁地扫描并清理无用数据。spark.cleaner.referenceTracking.cleanCheckpoints=true spark.cleaner.periodicGC.interval=30s - 分散临时文件存储路径:
将临时文件分散到多个磁盘分区,缓解单分区的空间压力。spark.local.dir=/disk1/spark-tmp,/disk2/spark-tmp
3. 拆分迭代任务为独立SparkSession
如果每次迭代的任务完全独立,可以为每个迭代创建单独的SparkSession,任务结束后关闭Session,Spark会自动清理所有关联的临时文件和缓存数据:
// 预先保存小型缓存DF到本地临时存储 filterDF.write.mode(SaveMode.Overwrite).parquet("/tmp/filterDF") noFilterDF.write.mode(SaveMode.Overwrite).parquet("/tmp/noFilterDF") fileList.sliding(500).foreach(x => { val session = SparkSession.builder().appName("SparkIterationTask").getOrCreate() // 加载预存的小型DF val filterDF = session.read.parquet("/tmp/filterDF") val noFilterDF = session.read.parquet("/tmp/noFilterDF") val data = session.read.format("avro").option("header", "true").load(x: _*) val joinedData = data.join(filterDF).join(noFilterDF) val results = joinedData.groupBy("x","y").agg(sum("p")) results.write.mode(SaveMode.Overwrite).parquet(filePath) // 关闭Session,自动清理所有临时数据 session.stop() })
这种方式彻底隔离了每次迭代的上下文,避免了跨迭代的Shuffle文件残留问题。
总结
第一次迭代的Shuffle数据未自动清理是同一SparkSession下的正常现象,因为Spark默认不会将循环内的任务视为独立上下文。手动删除磁盘文件会破坏BlockManager的元数据一致性,引发异常。通过显式清理临时DF、调整Spark配置或拆分独立Session的方式,可以有效解决磁盘空间被Shuffle文件填满的问题。
内容的提问来源于stack exchange,提问作者atul gurale

