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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 23:55:12