如何在SageMaker Notebook中将S3分区文件合并为单个输出文件?
解决方案:将S3上的PySpark分区文件合并为单个文件
针对超50GB的大型数据集,优先推荐基于PySpark的方案(避免内存瓶颈),以下是可行实现:
方案1:PySpark直接输出单个文件(最优)
既然已经用PySpark完成了数据合并,直接利用PySpark的分区控制将结果输出为单个文件,步骤如下:
代码实现
from pyspark.sql import SparkSession import boto3 # 初始化SparkSession(若未初始化) spark = SparkSession.builder.appName("MergeToSingleFile").getOrCreate() # 读取合并后的分区数据(替换为你的S3路径和文件格式) merged_df = spark.read.parquet("s3://your-bucket/path/to/merged-partitions/") # 合并为单个分区并写入临时目录(coalesce(1)避免shuffle,性能优于repartition(1)) merged_df.coalesce(1).write.mode("overwrite").parquet("s3://your-bucket/temp-single-file/") # 使用boto3将临时目录的part文件重命名为目标文件 s3_client = boto3.client("s3") bucket = "your-bucket" temp_prefix = "temp-single-file/" target_key = "final-merged-data.parquet" # 获取临时目录下的part文件 response = s3_client.list_objects_v2(Bucket=bucket, Prefix=temp_prefix) part_files = [obj["Key"] for obj in response["Contents"] if obj["Key"].endswith(".parquet")] if part_files: # 复制part文件到目标路径 s3_client.copy_object( Bucket=bucket, CopySource={"Bucket": bucket, "Key": part_files[0]}, Key=target_key ) # 清理临时文件和目录 for obj in response["Contents"]: s3_client.delete_object(Bucket=bucket, Key=obj["Key"]) s3_client.delete_object(Bucket=bucket, Key=temp_prefix)
注意:
- 若数据集极大,
coalesce(1)可能导致单个Executor内存压力过高,可先确保SageMaker Notebook的实例类型有足够内存(如ml.r5.2xlarge及以上)。 - 替换代码中的文件格式(如
csv、json)以匹配你的数据类型。
方案2:AWS CLI合并(仅适用于纯文本格式)
如果处理的是CSV/TXT等纯文本文件,可直接在SageMaker Notebook中通过AWS CLI合并:
命令示例
# 递归读取分区文件并合并为单个CSV aws s3 cp s3://your-bucket/path/to/merged-partitions/ s3://your-bucket/final-merged.csv \ --recursive --exclude "*" --include "*.csv" --concurrency 10
限制:该方法仅适用于纯文本格式,二进制格式(Parquet/ORC)不能直接合并,会损坏文件结构。
方案3:Pandas分块合并(备选,不推荐大文件)
若必须使用Pandas,可分块读取每个分区文件并追加到目标文件,避免一次性加载全量数据:
代码实现
import pandas as pd import boto3 s3_client = boto3.client("s3") bucket = "your-bucket" partition_prefix = "path/to/merged-partitions/" target_file = "s3://your-bucket/final-merged.csv" # 列出所有分区文件 response = s3_client.list_objects_v2(Bucket=bucket, Prefix=partition_prefix) file_keys = [obj["Key"] for obj in response["Contents"] if not obj["Key"].endswith("/")] first_write = True for key in file_keys: # 读取单个分区文件 obj = s3_client.get_object(Bucket=bucket, Key=key) chunk_df = pd.read_csv(obj["Body"]) # 替换为read_parquet等匹配格式的方法 # 追加写入目标文件 if first_write: chunk_df.to_csv(target_file, index=False) first_write = False else: chunk_df.to_csv(target_file, mode="a", header=False, index=False)
缺点:大文件下内存占用高、处理速度慢,仅适合小数据集或特殊场景。
内容的提问来源于stack exchange,提问作者Arcane
相关产品推荐
相关产品推荐

