【问题标题】:Write and read raw byte arrays in Spark - using Sequence File SequenceFile在 Spark 中写入和读取原始字节数组 - 使用序列文件 SequenceFile
【发布时间】:2014-07-27 19:30:01
【问题描述】:

如何使用 Apache Spark 将RDD[Array[Byte]] 写入文件并再次读回?

【问题讨论】:

    标签: scala hadoop hdfs apache-spark sequencefile


    【解决方案1】:

    常见问题似乎是一个奇怪的无法将异常从 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())
    

    【讨论】:

    • 这篇文章比较老所以只想知道答案是否仍然是最新的?阅读前还需要使用copyBytes吗?
    • @SamStoelinga 是的,我想是的,Hadoop API 不太可能改变。
    • 一个更有效的替代方法是使用<BytesWritableInstance>.getBytes() 并且只处理最多<BytesWritableInstance>.getLength() 字节。当然,如果你严格需要RDD[Array[Byte]],这种方法是行不通的,但你可以考虑RDD[(Array[Byte], Int)]
    • 任何人都可以发布完整的工作代码 sn-p 包括要导入的包吗?谢谢。
    • @Choix - 我遇到了同样的问题。发布解决了我的问题的 sn-p 作为单独的答案。
    【解决方案2】:

    这是一个 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)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多