如何使用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
相关产品推荐
相关产品推荐

