【发布时间】: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