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

如何将不同表头的多个CSV读取为单个Spark DataFrame?

问题描述

我有多个CSV文件,部分文件列名重叠,部分列名完全不同:

  • 文件1列:['circuitId', 'circuitRef', 'name', 'location', 'country', 'lat', 'lng', 'alt', 'url']
  • 文件2列:['raceId', 'year', 'round', 'circuitId', 'name', 'date', 'time', 'url']

想要创建包含所有列的DataFrame,尝试了预定义Schema的Spark读取代码,但输出不符合预期,代码如下:

from pyspark.sql.types import StructType, StructField, StringType

sch=StructType([StructField('circuitId',StringType(),True),
StructField('year',StringType(),True),
StructField('name',StringType(),True),
StructField('alt',StringType(),True),
StructField('url',StringType(),True),
StructField('round',StringType(),True),
StructField('lng',StringType(),True),
StructField('date',StringType(),True),
StructField('circuitRef',StringType(),True),
StructField('raceId',StringType(),True),
StructField('lat',StringType(),True),
StructField('location',StringType(),True),
StructField('country',StringType(),True),
StructField('time',StringType(),True)
])

df=spark.read \
        .option('header','true') \
        .schema(sch) \
        .csv('/FileStore/Udemy/Formula_One_Raw/*.csv')
问题原因及解决方案

直接指定全局Schema读取多CSV时,Spark会按Schema的列顺序匹配文件的列顺序,而非按列名匹配,这是输出不符合预期的核心原因——文件列顺序和Schema顺序不一致会导致数据映射错位。

方法1:自动合并所有列(推荐)

无需手动定义Schema,让Spark自动读取每个文件的Schema并合并,缺失列自动填充null:

# 读取所有CSV,自动合并Schema
df = spark.read \
    .option("header", "true") \
    .option("mergeSchema", "true") \  # 关键配置:合并所有文件的Schema
    .option("inferSchema", "true") \  # 可选,自动推断数据类型;关闭则所有列默认String类型
    .csv('/FileStore/Udemy/Formula_One_Raw/*.csv')
  • 优势:无需手动整理列,自动处理列名匹配与缺失值填充
  • 注意:数据量较大时,inferSchema会增加读取时间,可按需关闭

方法2:手动指定Schema并按列名匹配

若必须严格控制Schema,可逐个读取文件并补全缺失列后合并:

from pyspark.sql.functions import lit
import glob

# 复用你定义的完整Schema
sch=StructType([StructField('circuitId',StringType(),True),
StructField('year',StringType(),True),
StructField('name',StringType(),True),
StructField('alt',StringType(),True),
StructField('url',StringType(),True),
StructField('round',StringType(),True),
StructField('lng',StringType(),True),
StructField('date',StringType(),True),
StructField('circuitRef',StringType(),True),
StructField('raceId',StringType(),True),
StructField('lat',StringType(),True),
StructField('location',StringType(),True),
StructField('country',StringType(),True),
StructField('time',StringType(),True)
])

# 获取所有CSV文件路径
file_paths = glob.glob('/FileStore/Udemy/Formula_One_Raw/*.csv')

dfs = []
for path in file_paths:
    # 读取单个文件,按自身header匹配列
    temp_df = spark.read.option("header", "true").csv(path)
    # 补全Schema中存在但当前文件缺失的列,填充null
    for col_name in sch.names:
        if col_name not in temp_df.columns:
            temp_df = temp_df.withColumn(col_name, lit(None).cast(sch[col_name].dataType))
    # 按Schema列顺序调整
    temp_df = temp_df.select(sch.names)
    dfs.append(temp_df)

# 合并所有DataFrame
final_df = spark.createDataFrame(spark.sparkContext.emptyRDD(), sch)
for df in dfs:
    final_df = final_df.unionByName(df)
  • 优势:严格控制列类型与顺序
  • 劣势:需手动处理每个文件,代码量较大

内容的提问来源于stack exchange,提问作者Ankit Tyagi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 17:05:22