【问题标题】:How do I save a file in a Spark PairRDD using the key as the filename and the value as the contents?如何使用键作为文件名并使用值作为内容将文件保存在 Spark PairRDD 中?
【发布时间】:2016-04-05 21:14:25
【问题描述】:

在 Spark 中,我使用 sc.binaryFiles 从 s3 下载了多个文件。生成的 RDD 具有作为文件名的键,值具有文件的内容。我已经解压缩了文件内容,csv 对其进行了解析,并将其转换为数据帧。所以,现在我有一个 PairRDD[String, DataFrame]。我遇到的问题是我想使用密钥作为文件名将文件保存到 HDFS,并将值保存为镶木地板文件,如果它已经存在,则覆盖一个。这是我目前得到的。

val files = sc.binaryFiles(lFiles.mkString(","), 250).mapValues(stream => sc.parallelize(readZipStream(new ZipInputStream(stream.open))))
val tables = files.mapValues(file => {
    val header = file.first.split(",")
    val schema = StructType(header.map(fieldName => StructField(fieldName, StringType, true)))
    val lines = file.mapPartitionsWithIndex { (idx, iter) => if (idx == 0) iter.drop(1) else iter }.flatMap(x => x.split("\n"))
    val rowRDD = lines.map(x => Row.fromSeq(x.split(",")))
    sqlContext.createDataFrame(rowRDD, schema)
})

如果您有任何建议,请告诉我。我会很感激的。

谢谢, 本

【问题讨论】:

  • 天真的方法:如果您的密钥基数较低,您可以收集它们,对它们进行迭代过滤该密钥,然后将其写入磁盘,路径等于密钥。

标签: scala apache-spark rdd


【解决方案1】:

spark 中将文件保存到 HDFS 的方式与 hadoop 相同。所以你需要创建一个扩展MultipleTextOutputFormat的类,在自定义类中你可以自己定义输出文件名。示例如下:

class RDDMultipleTextOutputFormat extends MultipleTextOutputFormat[Any, Any] {
    override def generateFileNameForKeyValue(key: Any, value: Any, name: String): String = {
        "realtime-" + new SimpleDateFormat("yyyyMMddHHmm").format(new Date()) + "00-" + name
    }
}

调用代码如下:

RDD.rddToPairRDDFunctions(rdd.map { case (key, list) =>
    (NullWritable.get, key)
}).saveAsHadoopFile(input, classOf[NullWritable], classOf[String], classOf[RDDMultipleTextOutputFormat])

【讨论】:

  • 这真的可以在没有 HDFS 的情况下与 S3 Native FS 一起使用吗?我想知道文件何时会真正上传到 s3,可能是在工作完成时?因为最后 X 条记录可能属于所有 X 文件……所以在处理最后一条记录之前,什么都不能上传到 s3,对吧?
猜你喜欢
  • 2022-01-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-21
  • 2022-01-25
  • 2022-08-02
  • 2022-11-12
  • 2019-05-04
相关资源
最近更新 更多