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

如何读取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 14:15:22