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

PySpark:如何高效读取列位置不同的多个CSV文件

如何用Spark高效按列名合并多列不一致的CSV文件

直接使用Spark CSV数据源的mergeSchema参数,就能实现单次批量读取同时按列名正确合并数据,完美解决你遇到的效率与正确性矛盾问题。

最优方案代码

df = spark.read.format('csv') \
    .option('header', 'true') \
    .option('mergeSchema', 'true') \
    .load('/path/to/csv/folder/*.csv')

原理说明

mergeSchema=true会让Spark在读取前自动扫描所有目标CSV文件的表头,合并出一个包含所有出现过的列的统一Schema。读取文件时,每个文件的列会自动对应到统一Schema的对应位置,缺失的列会填充null,完全避免了按索引对齐导致的列错位问题。

这种方式保留了批量读取的高效性:所有文件由一个Spark作业并行处理,不需要像逐个读取后unionByName那样多次触发作业,性能和你测试的批量读取(6秒)基本一致。

进阶优化:指定自定义Schema

如果你提前知道所有可能的列名和对应数据类型,直接指定Schema可以跳过Spark自动合并Schema的扫描步骤,进一步提升效率:

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

# 根据实际业务定义完整的Schema
full_schema = StructType([
    StructField("A", StringType(), nullable=True),
    StructField("B", IntegerType(), nullable=True),
    StructField("C", StringType(), nullable=True)
])

df = spark.read.format('csv') \
    .option('header', 'true') \
    .schema(full_schema) \
    .load('/path/to/csv/folder/*.csv')

读取时Spark会直接按指定Schema对齐列,缺失列填充null,效率比mergeSchema更高。

对比原方案的优势

你之前用unionByName的方式慢,是因为每个文件都要单独触发一次读取任务,多次作业的调度和启动开销累积导致耗时增加。而批量读取+Schema合并是单作业处理,完全利用Spark的并行计算能力,性能提升明显。

内容的提问来源于stack exchange,提问作者Mr Pio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 19:16:12