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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 12:30:58