【发布时间】:2014-11-27 05:50:33
【问题描述】:
我有一些 spark 代码来处理 csv 文件。它对其进行了一些转换。我现在想将此 RDD 保存为 csv 文件并添加标题。此 RDD 的每一行都已正确格式化。
我不知道该怎么做。我想将标题字符串和我的 RDD 合并,但标题字符串不是 RDD,所以它不起作用。
【问题讨论】:
标签: apache-spark
我有一些 spark 代码来处理 csv 文件。它对其进行了一些转换。我现在想将此 RDD 保存为 csv 文件并添加标题。此 RDD 的每一行都已正确格式化。
我不知道该怎么做。我想将标题字符串和我的 RDD 合并,但标题字符串不是 RDD,所以它不起作用。
【问题讨论】:
标签: apache-spark
spark.sparkContext.parallelize(Seq(SqlHelper.getARow(temRet.columns,
temRet.columns.length))).union(temRet.rdd).map(x =>
x.mkString("\x01")).coalesce(1, true).saveAsTextFile(retPath)
object SqlHelper {
//create one row
def getARow(x: Array[String], size: Int): Row = {
var columnArray = new Array[String](size)
for (i <- 0 to (size - 1)) {
columnArray(i) = x(i).toString()
}
Row.fromSeq(columnArray)
}
}
【讨论】:
来自问题:我现在想将此 RDD 保存为 CSV 文件并添加标题。此 RDD 的每一行都已正确格式化。
使用 Spark 2.x,您可以通过多种方式选择convert RDD to DataFrame
val rdd = .... //Assume rdd properly formatted with case class or tuple
val df = spark.createDataFrame(rdd).toDF("col1", "col2", ... "coln")
df.write
.format("csv")
.option("header", "true") //adds header to file
.save("hdfs://location/to/save/csv")
现在我们甚至可以使用 Spark SQL DataFrame 来加载、转换和保存 CSV 文件
【讨论】:
def addHeaderToRdd(sparkCtx: SparkContext, lines: RDD[String], header: String): RDD[String] = {
val headerRDD = sparkCtx.parallelize(List((-1L, header))) // We index the header with -1, so that the sort will put it on top.
val pairRDD = lines.zipWithIndex()
val pairRDD2 = pairRDD.map(t => (t._2, t._1))
val allRDD = pairRDD2.union(headerRDD)
val allSortedRDD = allRDD.sortByKey()
return allSortedRDD.values
}
【讨论】:
在没有联合的情况下编写它的一些帮助(在合并时提供标题)
val fileHeader ="This is header"
val fileHeaderStream: InputStream = new ByteArrayInputStream(fileHeader.getBytes(StandardCharsets.UTF_8));
val output = IOUtils.copyBytes(fileHeaderStream,out,conf,false)
现在循环遍历文件部分以使用
编写完整的文件val in: DataInputStream = ...<data input stream from file >
IOUtils.copyBytes(in, output, conf, false)
这对我来说确保了标题始终位于第一行,即使您使用 coalasec/repartition 进行高效写入
【讨论】:
您可以从标题行中创建一个 RDD,然后 union 它,是的:
val rdd: RDD[String] = ...
val header: RDD[String] = sc.parallelize(Array("my,header,row"))
header.union(rdd).saveAsTextFile(...)
然后你会得到一堆你合并的part-xxxxx文件。
问题是我认为你不能保证标题将是第一个分区,因此最终会出现在part-00000 和文件的顶部。在实践中,我很确定它会。
更可靠的是使用像hdfs 这样的Hadoop 命令来合并part-xxxxx 文件,并且作为命令的一部分,只需从文件中放入标题行。
【讨论】:
val header = sc.parallelize(Array('col1','col2'), 1) header.union( rdd.map(_.toString)) .repartition(1).saveAsTextFile(outputLocation)
union() 不保证保持秩序。现在正在寻找解决方法,看起来对 RDD 进行排序可能会有所帮助