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

Spark Scala中No space left on device错误的排查与解决

Spark迭代任务中Shuffle临时文件未自动清理问题解析与解决

问题根源解析

  1. Shuffle文件的管理逻辑
    Spark的Shuffle临时文件由Executor的BlockManager负责管理,只有当关联的RDD/DataSet的所有依赖计算完成,且Spark确认这些数据不会被再次复用的情况下,才会触发自动清理。在你的循环迭代场景中,所有任务共享同一个SparkSession,Spark不会将每次迭代视为独立的上下文,因此第一次迭代生成的Shuffle文件会被BlockManager保留,不会自动清理——这不属于默认的自动清理场景。

  2. 手动删除引发异常的原因
    你直接删除/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服务:
    spark.shuffle.service.enabled=true
    
    外部Shuffle服务会独立管理Shuffle文件,避免Executor退出时的文件残留,同时更智能地清理不再使用的Shuffle数据。
  • 调整Shuffle文件清理阈值:
    spark.cleaner.referenceTracking.cleanCheckpoints=true
    spark.cleaner.periodicGC.interval=30s
    
    开启检查点清理并缩短GC跟踪的间隔,让ContextCleaner更频繁地扫描并清理无用数据。
  • 分散临时文件存储路径:
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:45:01