Spark 2.1:从空Dataset<Row>生成带表头CSV并规避空文件
解决Spark 2.1中空Dataset生成带表头CSV的问题
在Spark 2.1版本里,确实存在空Dataset写入CSV时只生成空文件、不输出表头的问题——因为Spark此时会判断没有数据,直接跳过了表头的写入逻辑。下面给你两种可行的实现方案:
方案一:手动判断Dataset状态,单独生成表头文件
这种方法最直接,先检查Dataset是否为空,如果是空的,就手动输出表头;如果有数据,正常写入即可。
Scala 示例代码
import org.apache.spark.sql.{Dataset, Row, SparkSession} import org.apache.hadoop.fs.{FileSystem, Path} import java.io.OutputStreamWriter def writeCsvWithHeader(ds: Dataset[Row], outputPath: String)(implicit spark: SparkSession): Unit = { if (ds.isEmpty) { // 提取表头字段,拼接成CSV格式的表头字符串 val header = ds.schema.fields.map(_.name).mkString(",") // 针对HDFS/本地文件系统写入表头文件 val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration) val outputStream = fs.create(new Path(s"$outputPath/part-00000.csv")) val writer = new OutputStreamWriter(outputStream) writer.write(header) writer.close() outputStream.close() // 生成Spark标准输出的_SUCCESS标记文件 fs.create(new Path(s"$outputPath/_SUCCESS")).close() } else { ds.write .option("header", "true") .csv(outputPath) } }
Java 示例代码
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import java.io.OutputStreamWriter; import java.util.Arrays; public class CsvWriterHelper { public static void writeCsvWithHeader(Dataset<Row> ds, String outputPath, SparkSession spark) { if (ds.isEmpty()) { // 构建表头字符串 String[] columnNames = ds.schema().fieldNames(); String header = String.join(",", columnNames); try { FileSystem fs = FileSystem.get(spark.sparkContext().hadoopConfiguration()); // 写入表头文件 OutputStreamWriter writer = new OutputStreamWriter(fs.create(new Path(outputPath + "/part-00000.csv"))); writer.write(header); writer.close(); // 生成_SUCCESS标记 fs.create(new Path(outputPath + "/_SUCCESS")).close(); } catch (Exception e) { e.printStackTrace(); } } else { ds.write() .option("header", "true") .csv(outputPath); } } }
方案二:Union一个空行(仅适用于可接受空数据行的场景)
如果可以容忍输出文件里除了表头还有一行空数据,也可以给空Dataset union一条符合schema的空Row,让Spark认为存在数据,从而触发表头写入:
Scala 示例
import org.apache.spark.sql.Row val emptyRow = Row.fromSeq(ds.schema.fields.map(_ => null)) val dummyDs = spark.createDataFrame(spark.sparkContext.parallelize(Seq(emptyRow)), ds.schema) ds.union(dummyDs) .write .option("header", "true") .csv(outputPath)
这种方法的缺点是会多一行空数据,如果你需要纯表头的CSV文件,方案一更合适。
内容的提问来源于stack exchange,提问作者Whimsical
相关产品推荐
相关产品推荐

