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

如何使用PySpark将文件夹中的所有CSV文件加载为单个DataFrame

使用PySpark合并文件夹中的多个CSV文件为单个DataFrame

核心方法:直接读取文件夹路径

PySpark的read.csvAPI支持直接读取整个文件夹路径,自动加载文件夹内所有CSV文件并合并为单一DataFrame,前提是文件结构(列、数据类型)一致。

基础代码示例

from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder \
    .appName("MergeSalesCSVs") \
    .getOrCreate()

# 读取目标文件夹下所有CSV,合并为DataFrame
merged_df = spark.read.csv(
    "/path/to/your/sales-folder",
    header=True,  # 指明CSV首行为列名
    inferSchema=True  # 自动推断列的数据类型,数据量较大时建议手动指定schema
)

# 验证结果
merged_df.printSchema()
merged_df.show(5)  # 显示前5行数据

处理结构不一致的CSV文件

如果文件夹内的CSV存在列数不同、列名差异的情况,添加mergeSchema=True参数,Spark会自动合并所有文件的Schema,缺失的列将填充null:

merged_df = spark.read.csv(
    "/path/to/your/sales-folder",
    header=True,
    inferSchema=True,
    mergeSchema=True
)

读取特定前缀的CSV文件

如果文件夹内还有其他无关CSV,可通过通配符指定只读取Sales_开头的文件:

merged_df = spark.read.csv(
    "/path/to/your/sales-folder/Sales_*.csv",
    header=True,
    inferSchema=True
)

优化:手动指定Schema

当数据量较大时,inferSchema会额外扫描数据,影响效率,建议手动定义Schema:

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

# 自定义与CSV匹配的Schema
sales_schema = StructType([
    StructField("TransactionDate", StringType(), nullable=True),
    StructField("ProductName", StringType(), nullable=True),
    StructField("SalesAmount", DoubleType(), nullable=True),
    StructField("QuantitySold", IntegerType(), nullable=True)
])

# 使用自定义Schema读取
merged_df = spark.read.csv(
    "/path/to/your/sales-folder",
    header=True,
    schema=sales_schema
)

注意事项

  • 确保文件夹路径正确,若为本地路径需保证Spark有权限访问;若为HDFS路径,需写成hdfs:///path/to/folder格式;
  • 若CSV使用非逗号分隔符(如分号),可通过sep=";"参数指定分隔符。

内容的提问来源于stack exchange,提问作者Sahand Pourjavad

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 10:10:23