如何在PySpark中将S3指定文件夹下的多个CSV合并为单个Parquet文件
使用PySpark合并多个CSV为单个Parquet文件
核心思路
直接读取目标S3目录下的所有CSV文件(无需逐个指定文件名),合并后写入单个Parquet文件,通过coalesce(1)控制输出文件数量。
完整代码实现
from pyspark.sql import SparkSession # 初始化Spark会话 spark = SparkSession.builder \ .appName("MergeCSVToParquet") \ .getOrCreate() # 读取指定S3目录下的全部CSV文件 source_csv_dir = "s3://lla.raw.dev/data/shared/sap/orders/2022/09/10/" orders_df = spark.read.csv( source_csv_dir, header=True, # 若CSV无表头,改为False并手动定义schema inferSchema=True # 自动推断字段类型,生产环境建议手动指定schema以保证稳定性 ) # 将合并后的数据写入单个Parquet文件 target_parquet_path = "s3://lla.raw.dev/data/shared/sap/orders/merged_orders.parquet" orders_df.coalesce(1) \ .write \ .mode("overwrite") \ # 可选模式:append/ignore/error .parquet(target_parquet_path) # 关闭Spark会话 spark.stop()
关键注意点
coalesce(1)的使用:该操作将数据合并到一个分区,从而生成单个Parquet文件,但仅适合数据量较小的场景。如果数据量庞大,强行合并会导致单节点压力过大,建议根据实际情况调整分区数,或放弃单个文件的需求。- Schema处理:生产环境中建议手动定义schema(而非
inferSchema=True),避免自动推断出错,示例如下:from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType custom_schema = StructType([ StructField("order_id", IntegerType(), nullable=False), StructField("customer_id", StringType(), nullable=True), StructField("order_date", DateType(), nullable=True), # 根据你的CSV字段补充其他列定义 ]) orders_df = spark.read.csv(source_csv_dir, header=True, schema=custom_schema) - 写入模式:
mode("overwrite")会覆盖目标路径下的已有文件,若需要追加数据可改为mode("append"),默认模式为error(若目标路径存在则抛出异常)。 - S3权限:确保Spark集群具备访问目标S3路径的读写权限,必要时配置AWS认证参数(如
spark.hadoop.fs.s3a.access.key等)。
内容的提问来源于stack exchange,提问作者Anos
相关产品推荐
相关产品推荐

