如何用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
相关产品推荐
相关产品推荐

