如何读取S3中Schema不同的CSV文件并合并为完整DataFrame?
问题:合并读取S3上列结构不同的CSV文件并保留所有列
我在S3上有两个CSV文件:
a1.csv内容:
# a1.csv a,b 3,4
b2.csv内容:
# b2.csv a,c 1,"text"
期望读取后得到包含所有列的DataFrame,结果如下:
+---+----+----+ | a| b| c| +---+----+----+ | 1|null|text| | 3| 4|null| +---+----+----+
尝试两种方式均未得到预期结果:
方案1:使用inferSchema参数
df = spark.read\ .option("header", True)\ .option("inferSchema", True)\ .csv("s3a://test/*.csv")\ .show()
输出结果:
+---+----+ | a| c| +---+----+ | 1|text| | 3| 4| +---+----+
方案2:指定自定义Schema
from pyspark.sql.types import StructType, StructField, IntegerType, StringType schema = StructType([ StructField("a", IntegerType(), False) ,StructField("b", IntegerType(), True) ,StructField("c", StringType(), True) ]) df = spark.read\ .option("header", True)\ .schema(schema)\ .csv("s3a://test/*.csv")\ .show()
输出结果:
+---+----+----+ | a| b| c| +---+----+----+ | 1|null|null| | 3| 4|null| +---+----+----+
请问有什么方法可以实现需求?
解决方案
直接用通配符读取时,Spark会以首个扫描到的文件列结构为基准,或忽略Schema中未匹配的列值。要保留所有列,需分别读取单个文件后用unionByName合并:
# 读取单个文件 df_a = spark.read.option("header", True).option("inferSchema", True).csv("s3a://test/a1.csv") df_b = spark.read.option("header", True).option("inferSchema", True).csv("s3a://test/b2.csv") # 按列名合并,缺失列自动填充null final_df = df_a.unionByName(df_b, allowMissingColumns=True) final_df.show()
执行后输出符合预期:
+---+----+----+ | a| b| c| +---+----+----+ | 3| 4|null| | 1|null|text| +---+----+----+
补充说明
unionByName会自动对齐两个DataFrame的列名,allowMissingColumns=True允许合并时存在缺失列,对应位置填充null。- 若文件数量较多,可遍历S3目录下的文件路径,逐个读取后依次执行合并操作。
内容的提问来源于stack exchange,提问作者ohh
相关产品推荐
相关产品推荐

