【发布时间】:2016-09-03 16:30:06
【问题描述】:
我需要在每个 map() 中读取不同的文件,该文件在 HDFS 中
val rdd=sc.parallelize(1 to 10000)
val rdd2=rdd.map{x=>
val hdfs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI("hdfs://ITS-Hadoop10:9000/"), new org.apache.hadoop.conf.Configuration())
val path=new Path("/user/zhc/"+x+"/")
val t=hdfs.listStatus(path)
val in =hdfs.open(t(0).getPath)
val reader = new BufferedReader(new InputStreamReader(in))
var l=reader.readLine()
}
rdd2.count
我的问题是这段代码
val hdfs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI("hdfs://ITS-Hadoop10:9000/"), new org.apache.hadoop.conf.Configuration())
运行时间太长,每次 map() 都需要创建一个新的 FileSystem 值。我可以把这段代码放在 map() 函数之外,这样它就不必每次都创建 hdfs 了吗?或者如何在 map() 中快速读取文件?
我的代码在多台机器上运行。谢谢!
【问题讨论】:
-
尝试将
val hdfs移出地图封闭区。 -
-
@tuxdna 我试图将它放在地图关闭之外,但它有错误“任务不可序列化,由:java.io.NotSerializableException:org.apache.hadoop.hdfs.DistributedFileSystem 引起”
-
@eliasah 文件很小,但我不完全理解您所说的内容,您是否建议将我需要的所有文件加载到类似于 Knows Not Much 建议的 RDD 中?
-
老实说,我不明白 KnowsNotMuch 建议什么。
标签: scala apache-spark