Scala+Spark合并多个CSV文件时表头重复或缺失问题
Spark合并多CSV文件单表头输出实现方案
问题说明
需要合并多个表头结构完全一致、数据内容不同的CSV文件为单个文件,待合并文件按data_0_1、data_0_2规则依次编号。
原有实现采用Spark+Scala编写,代码如下:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ import org.apache.spark.sql.{Dataset, Row} import spark.implicits._ val INPUT_BUCKET_PREFIX = "fie:/path/data/"; def getData(tableName: String): Dataset[Row] = { spark.read .option("header", "true") .option("ignoreLeadingWhiteSpace", "true") .option("ignoreTrailingWhiteSpace", "true") .csv(INPUT_BUCKET_PREFIX + tableName) } getData("data*") .coalesce(1) .write.csv("file:/path/output")
当前存在的异常:
- 读取配置
header=true时,输出文件会重复多次写入表头,不符合要求 - 写入时不配置
header=true,输出文件完全不生成表头 - 目标效果:输出文件仅在首行写入一次表头
原因分析
出现重复表头的核心原因有两个:
- 原有代码写入阶段没有显式配置
header=true,Spark写入CSV时默认不输出表头,读取阶段的header配置和写入阶段相互独立,不会自动传递 - 代码中存在路径笔误(
fie:/应为file:/),同时未强制schema校验,部分文件的表头可能因为格式差异未被Spark识别为元数据,被当做普通数据行读入,最终写入输出文件造成重复
实现方案
方案一:常规场景最简修正(适合所有文件表头完全规范一致的场景)
直接修正原有代码的笔误,添加强制schema校验和写入阶段的header配置即可,代码如下:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ import org.apache.spark.sql.{Dataset, Row} import spark.implicits._ // 修正路径笔误 val INPUT_BUCKET_PREFIX = "file:/path/data/" def getData(tableName: String): Dataset[Row] = { spark.read .option("header", "true") .option("ignoreLeadingWhiteSpace", "true") .option("ignoreTrailingWhiteSpace", "true") // 强制按统一schema解析,避免表头错位 .option("enforceSchema", "true") .csv(INPUT_BUCKET_PREFIX + tableName) } getData("data*") .coalesce(1) .write // 写入阶段显式开启表头输出,仅会在输出文件首行写入一次表头 .option("header", "true") .csv("file:/path/output")
方案二:稳妥兼容方案(适合存在个别文件表头格式不规范、方案一仍出现重复表头的场景)
先读取首个文件获取标准表头,再全局读取所有文件时手动过滤掉各文件的表头行,从根源避免表头被当做数据写入:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ import org.apache.spark.sql.{Dataset, Row} import spark.implicits._ val INPUT_BUCKET_PREFIX = "file:/path/data/" // 读取第一个文件获取标准表头和schema val firstFileDf = spark.read .option("header", "true") .option("ignoreLeadingWhiteSpace", "true") .option("ignoreTrailingWhiteSpace", "true") .csv(INPUT_BUCKET_PREFIX + "data_0_1") val standardSchema = firstFileDf.schema val standardHeaderLine = firstFileDf.columns.mkString(",").trim // 全局读取所有文件,关闭自动表头识别,按标准schema解析 val allRawDf = spark.read .option("header", "false") .option("ignoreLeadingWhiteSpace", "true") .option("ignoreTrailingWhiteSpace", "true") .schema(standardSchema) .csv(INPUT_BUCKET_PREFIX + "data*") // 过滤掉所有和表头内容一致的行(即每个文件自带的表头行) val cleanDataDf = allRawDf.filter(row => row.mkString(",").trim != standardHeaderLine) // 合并为单分区写入,开启表头输出 cleanDataDf.coalesce(1) .write .option("header", "true") .csv("file:/path/output")
补充说明
coalesce(1)会将所有数据合并到1个分区,最终输出目录中只会生成1个数据分片,配合写入阶段的header=true配置,只会输出一次表头- Spark写入的输出目录中,除了命名格式为
part-00000-xxxx.csv的结果文件,还会生成_SUCCESS标记文件和隐藏的校验文件,需要最终CSV文件的话,直接将part开头的csv文件重命名移出即可
内容的提问来源于stack exchange,提问作者Divakar R
相关产品推荐
相关产品推荐

