如何将不同表头的多个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
相关产品推荐
相关产品推荐

