如何用Spark合并多子目录文件并保留原目录结构?
Spark合并S3各子目录下CSV文件为单个merged.csv
核心思路
不需要手动启动多个worker遍历目录,通过以下两步即可实现需求:
- 批量筛选S3存储桶下包含CSV文件的子目录
- 对每个目标子目录,读取所有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
相关产品推荐
相关产品推荐

