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

如何用Spark合并多子目录文件并保留原目录结构?

Spark合并S3各子目录下CSV文件为单个merged.csv

核心思路

不需要手动启动多个worker遍历目录,通过以下两步即可实现需求:

  1. 批量筛选S3存储桶下包含CSV文件的子目录
  2. 对每个目标子目录,读取所有CSV文件合并为单个文件,替换原目录下的分块文件

实现方案(Python版)

1. 初始化SparkSession与S3客户端

from pyspark.sql import SparkSession
import boto3

# 初始化SparkSession,按需配置S3权限参数
spark = SparkSession.builder \
    .appName("S3CSVMerger") \
    .config("spark.hadoop.fs.s3a.access.key", "your-access-key") \
    .config("spark.hadoop.fs.s3a.secret.key", "your-secret-key") \
    .getOrCreate()

# 初始化S3客户端用于目录遍历和文件操作
s3_client = boto3.client('s3')
bucket_name = "your-s3-bucket-name"

2. 筛选出包含CSV文件的子目录

# 获取存储桶下所有一级子目录
paginator = s3_client.get_paginator('list_objects_v2')
response_iterator = paginator.paginate(Bucket=bucket_name, Delimiter='/')

target_dirs = []
for response in response_iterator:
    if 'CommonPrefixes' in response:
        for prefix in response['CommonPrefixes']:
            subdir = prefix['Prefix']
            # 检查该目录下是否存在CSV文件
            obj_list = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=subdir)
            if 'Contents' in obj_list:
                has_csv = any(obj['Key'].endswith('.csv') for obj in obj_list['Contents'])
                if has_csv:
                    target_dirs.append(subdir)

3. 遍历目录合并CSV文件

for dir_path in target_dirs:
    # 读取当前目录下所有CSV文件
    input_path = f"s3://{bucket_name}/{dir_path}*.csv"
    df = spark.read.csv(input_path, header=True, inferSchema=True)
    
    # 合并为单个分区(coalesce避免不必要的shuffle,比repartition性能更优)
    merged_df = df.coalesce(1)
    
    # 写入临时目录(Spark写入会生成part-*文件,需后续重命名)
    temp_output = f"s3://{bucket_name}/{dir_path}temp_merge"
    merged_df.write.csv(temp_output, header=True, mode="overwrite")
    
    # 定位临时目录下的CSV文件,重命名为merged.csv
    temp_objs = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=f"{dir_path}temp_merge/")
    part_file = None
    for obj in temp_objs['Contents']:
        if obj['Key'].endswith('.csv') and not obj['Key'].endswith('_SUCCESS'):
            part_file = obj['Key']
            break
    
    if part_file:
        # 复制并改名到目标路径
        s3_client.copy_object(
            Bucket=bucket_name,
            CopySource={'Bucket': bucket_name, 'Key': part_file},
            Key=f"{dir_path}merged.csv"
        )
        # 删除临时目录文件
        for obj in temp_objs['Contents']:
            s3_client.delete_object(Bucket=bucket_name, Key=obj['Key'])
        # 删除原目录下的part_*.csv文件
        original_objs = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=dir_path)
        for obj in original_objs['Contents']:
            if obj['Key'].startswith(f"{dir_path}part_") and obj['Key'].endswith('.csv'):
                s3_client.delete_object(Bucket=bucket_name, Key=obj['Key'])

# 关闭SparkSession
spark.stop()

优化与注意事项

  • 并行处理:如果子目录数量较多,可将目录列表转为RDD并行处理,提升效率:
    def process_single_dir(dir_path):
        # 封装单个目录的处理逻辑
        input_path = f"s3://{bucket_name}/{dir_path}*.csv"
        df = spark.read.csv(input_path, header=True, inferSchema=True)
        merged_df = df.coalesce(1)
        temp_output = f"s3://{bucket_name}/{dir_path}temp_merge"
        merged_df.write.csv(temp_output, header=True, mode="overwrite")
        
        # 后续文件操作逻辑(需确保worker节点配置S3权限)
    
    # 并行处理所有目录
    spark.sparkContext.parallelize(target_dirs).foreach(process_single_dir)
    
  • 性能优化:优先使用coalesce(1)而非repartition(1),前者不会触发数据shuffle,大幅减少资源消耗
  • 权限配置:确保Spark集群和S3客户端拥有目标存储桶的读写权限
  • 空目录跳过:代码已自动筛选出包含CSV文件的目录,空目录会被忽略
  • 大文件处理:如果单目录下数据量极大,coalesce(1)可能导致单个executor内存不足,需根据实际情况调整或拆分处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 14:01:17