【发布时间】:2014-07-27 19:30:01
【问题描述】:
如何使用 Apache Spark 将RDD[Array[Byte]] 写入文件并再次读回?
【问题讨论】:
标签: scala hadoop hdfs apache-spark sequencefile
如何使用 Apache Spark 将RDD[Array[Byte]] 写入文件并再次读回?
【问题讨论】:
标签: scala hadoop hdfs apache-spark sequencefile
常见问题似乎是一个奇怪的无法将异常从 BytesWritable 转换为 NullWritable。其他常见问题是 BytesWritable getBytes 是一堆毫无意义的废话,根本没有得到字节。 getBytes 所做的是获取您的字节,而不是在最后添加大量零!你必须使用copyBytes
val rdd: RDD[Array[Byte]] = ???
// To write
rdd.map(bytesArray => (NullWritable.get(), new BytesWritable(bytesArray)))
.saveAsSequenceFile("/output/path", codecOpt)
// To read
val rdd: RDD[Array[Byte]] = sc.sequenceFile[NullWritable, BytesWritable]("/input/path")
.map(_._2.copyBytes())
【讨论】:
<BytesWritableInstance>.getBytes() 并且只处理最多<BytesWritableInstance>.getLength() 字节。当然,如果你严格需要RDD[Array[Byte]],这种方法是行不通的,但你可以考虑RDD[(Array[Byte], Int)]。
这是一个 sn-p,其中包含所有必需的导入,您可以根据 @Choix 的要求从 spark-shell 运行它们
import org.apache.hadoop.io.BytesWritable
import org.apache.hadoop.io.NullWritable
val path = "/tmp/path"
val rdd = sc.parallelize(List("foo"))
val bytesRdd = rdd.map{str => (NullWritable.get, new BytesWritable(str.getBytes) ) }
bytesRdd.saveAsSequenceFile(path)
val recovered = sc.sequenceFile[NullWritable, BytesWritable]("/tmp/path").map(_._2.copyBytes())
val recoveredAsString = recovered.map( new String(_) )
recoveredAsString.collect()
// result is: Array[String] = Array(foo)
【讨论】: