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

如何用AWS Glue将S3多目录JSON合并为指定大小的Parquet文件

实现S3多目录JSON合并为指定大小的Parquet文件

针对你的需求,这里提供三种可行方案,分别适配不同场景:

一、PySpark方案(推荐用于大规模数据)

PySpark天生适合处理分布式存储(如S3)的批量数据,能自动按指定大小拆分输出文件,操作高效且稳定。

步骤与代码示例

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# 初始化SparkSession,设置单个Parquet文件最大为100MB(单位:字节)
spark = SparkSession.builder \
    .appName("JSONToParquetWithSizeLimit") \
    .config("spark.sql.files.maxPartitionBytes", "104857600")  # 100*1024*1024
    .getOrCreate()

# 【推荐】提前定义JSON的Schema,避免自动推断字段时出错
# 根据你的实际JSON字段结构调整以下内容
custom_schema = StructType([
    StructField("user_id", IntegerType(), nullable=True),
    StructField("username", StringType(), nullable=True),
    StructField("create_time", StringType(), nullable=True)
])

# 递归读取S3所有子目录下的JSON文件
# 路径中的**表示遍历所有层级的子目录
df = spark.read.schema(custom_schema).json("s3://your-input-bucket/**/*.json")

# 写出Parquet文件到S3,Spark会自动按设置的100MB大小拆分文件
# mode("overwrite")表示覆盖已有输出,根据需要可改为"append"
df.write.mode("overwrite").parquet("s3://your-output-bucket/parquet-results/")

spark.stop()

关键说明

  • spark.sql.files.maxPartitionBytes:控制每个数据分区的大小,每个分区对应一个输出Parquet文件,因此设置为100MB后,输出文件会尽量接近该大小(如你的示例中生成96MB和64MB的文件)。
  • 提前定义Schema:即使你确认所有JSON字段一致,显式定义Schema也能避免Spark因个别文件的异常格式导致的推断错误,提升稳定性。
  • S3路径通配符:**会递归遍历所有子目录,确保读取到所有目标JSON文件。

二、AWS Glue方案(AWS生态用户首选)

AWS Glue是托管式Spark服务,无需自行维护集群,适合AWS环境下的批量数据处理,核心逻辑与PySpark一致。

步骤与代码示例

import sys
from awsglue.context import GlueContext
from awsglue.job import Job
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# 初始化Glue上下文,设置文件大小限制
glueContext = GlueContext(SparkSession.builder.config("spark.sql.files.maxPartitionBytes", "104857600").getOrCreate())
job = Job(glueContext)
job.init(sys.argv[1], sys.argv[2])

# 定义Schema
custom_schema = StructType([
    StructField("user_id", IntegerType(), nullable=True),
    StructField("username", StringType(), nullable=True),
    StructField("create_time", StringType(), nullable=True)
])

# 读取S3上的JSON文件(递归遍历子目录)
dynamic_frame = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": ["s3://your-input-bucket/"], "recurse": True},
    format="json",
    format_options={"schema": custom_schema.json()}
)

# 转换为Spark DataFrame
df = dynamic_frame.toDF()

# 写出Parquet到S3
df.write.mode("overwrite").parquet("s3://your-output-bucket/glue-parquet-results/")

job.commit()

关键说明

  • 配置IAM角色:确保Glue Job的IAM角色拥有S3读写权限,否则无法访问输入输出路径。
  • 无需集群管理:Glue会自动分配计算资源,完成后释放,适合周期性的批量处理任务。

三、Pandas+PyArrow方案(适合小数据量场景)

如果你的数据量不大(如示例中的20个8MB文件),可以用Pandas结合PyArrow手动分批合并,灵活性更高,但要注意内存占用。

步骤与代码示例

import boto3
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
import os

# 配置参数
input_bucket = "your-input-bucket"
output_bucket = "your-output-bucket"
output_prefix = "parquet-output/"
max_file_size_mb = 100
max_file_size_bytes = max_file_size_mb * 1024 * 1024

# 初始化S3客户端
s3 = boto3.client("s3")

# 递归列出S3上所有JSON文件
def get_all_json_files(bucket, prefix=""):
    paginator = s3.get_paginator("list_objects_v2")
    json_files = []
    for page in paginator.paginate(Bucket=bucket, Prefix=prefix):
        for obj in page.get("Contents", []):
            if obj["Key"].endswith(".json"):
                json_files.append(obj["Key"])
    return json_files

json_files = get_all_json_files(input_bucket)

# 分批合并逻辑
current_batch = []
current_estimated_size = 0
batch_index = 1

for file_key in json_files:
    # 临时下载JSON文件到本地
    temp_path = f"/tmp/{os.path.basename(file_key)}"
    s3.download_file(input_bucket, file_key, temp_path)
    
    # 读取JSON(如果是行分隔JSON,添加lines=True)
    df = pd.read_json(temp_path)
    os.remove(temp_path)
    
    # 估算当前DataFrame转Parquet后的大小(近似值)
    table = pa.Table.from_pandas(df)
    estimated_size = table.nbytes
    
    # 检查是否超过大小限制
    if current_estimated_size + estimated_size > max_file_size_bytes and current_batch:
        # 合并当前批次并保存
        batch_df = pd.concat(current_batch, ignore_index=True)
        output_file = f"batch_{batch_index}.parquet"
        local_output_path = f"/tmp/{output_file}"
        
        # 保存为Parquet
        batch_df.to_parquet(local_output_path, engine="pyarrow")
        
        # 上传到S3
        s3.upload_file(local_output_path, output_bucket, f"{output_prefix}{output_file}")
        os.remove(local_output_path)
        
        # 重置批次
        current_batch = [df]
        current_estimated_size = estimated_size
        batch_index += 1
    else:
        current_batch.append(df)
        current_estimated_size += estimated_size

# 处理最后一批剩余数据
if current_batch:
    batch_df = pd.concat(current_batch, ignore_index=True)
    output_file = f"batch_{batch_index}.parquet"
    local_output_path = f"/tmp/{output_file}"
    batch_df.to_parquet(local_output_path, engine="pyarrow")
    s3.upload_file(local_output_path, output_bucket, f"{output_prefix}{output_file}")
    os.remove(local_output_path)

关键说明

  • 大小估算为近似值:Parquet的压缩率会影响实际文件大小,因此最终输出可能略有偏差,但整体符合100MB的限制要求。
  • 内存注意事项:如果单批数据过大,可能导致内存溢出,因此该方案仅适合数据量较小的场景。

内容的提问来源于stack exchange,提问作者D.E. Veloper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 16:05:28