You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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,输出文件完全不生成表头
  • 目标效果:输出文件仅在首行写入一次表头

原因分析

出现重复表头的核心原因有两个:

  1. 原有代码写入阶段没有显式配置header=true,Spark写入CSV时默认不输出表头,读取阶段的header配置和写入阶段相互独立,不会自动传递
  2. 代码中存在路径笔误(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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 21:54:29