Spark集群循环删除HDFS文件时抛出NotSerializableException异常
解决Spark中删除HDFS文件时的SerializableException问题
嗨,这个坑我踩过好多次了!你遇到的java.io.NotSerializableException: org.apache.hadoop.hdfs.DistributedFileSystem异常,根源在于Hadoop的FileSystem实例是不可序列化的客户端对象,而你的代码写法不小心让Spark尝试把它序列化后传到Executor节点,自然就报错了。另外还要提一句:你原来的代码漏了iteration += 1,会导致无限循环,这个小细节可不能忘!
问题出在哪?
你在循环的if块里每次都创建FileSystem实例,这块代码哪怕是在Driver端跑循环,也可能被Spark的闭包机制捕获,触发序列化逻辑——毕竟FileSystem本身没有实现Serializable接口,根本没法被序列化传输。
靠谱的解决方案
方案1:提前在Driver端初始化FileSystem(最简单)
把FileSystem的获取移到循环外面,只初始化一次,确保它全程在Driver端运行,不会被序列化到Executor:
val numIter=10 var iteration = 0 // 提前在Driver端获取FileSystem实例,只做一次初始化 val fs = FileSystem.get(sc.hadoopConfiguration) while (iteration < numIter) { val b = iteration - 2 if(b >= 0){ val edgpath = "Mypath" + b fs.delete(new Path(edgpath), true) } graph.vertices.saveAsObjectFile("Mypath_" + iteration) iteration += 1 // 必须加这句,否则循环会无限执行! }
方案2:用Spark官方工具类(更稳妥)
Spark提供了SparkHadoopUtil工具类,专门用来处理Hadoop相关的配置和对象初始化,能自动避开序列化陷阱:
import org.apache.spark.deploy.SparkHadoopUtil import org.apache.hadoop.fs.Path val numIter=10 var iteration = 0 while (iteration < numIter) { val b = iteration - 2 if(b >= 0){ val edgpath = "Mypath" + b // 用SparkHadoopUtil获取安全的Hadoop配置和FileSystem val hadoopConf = SparkHadoopUtil.get.newConfiguration(sc.hadoopConfiguration) FileSystem.get(hadoopConf).delete(new Path(edgpath), true) } graph.vertices.saveAsObjectFile("Mypath_" + iteration) iteration += 1 }
方案3:明确在Driver端执行删除操作
如果你确定整个循环是在Driver端运行的(大部分场景下都是),只要确保删除文件的逻辑不被包含在RDD的算子闭包里(比如map、foreach这些会发到Executor的操作),就不会有问题。方案1其实就是这个思路,把FileSystem固定在Driver端使用。
避坑小贴士
- 绝对不要在RDD的
map/flatMap/foreach等算子内部创建FileSystem实例!这些代码会在Executor节点上执行,重复创建客户端不仅浪费资源,还容易引发连接或序列化问题。 - 如果必须在Executor上操作HDFS,优先用Hadoop的
Path结合配置,或者通过Spark的API间接实现,别直接碰FileSystem实例。
内容的提问来源于stack exchange,提问作者Ashwini
相关产品推荐
相关产品推荐

