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
相关产品推荐
相关产品推荐

