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

如何优化处理S3海量小JSON文件的AWS Glue PySpark作业速度?

问题描述

有一个用于为Adobe Experience Platform加载历史数据的AWS Glue作业,需处理S3中10万+个小JSON文件(每个文件5-50KB,含单条JSON记录,总数据量约5GB)。当前作业基于AWS Glue 3.0(Spark 3.1),采用Spark直接读取JSON文件的方式,耗时约4小时,无法满足业务需求。已尝试增加DPU、调整spark.sql.files.maxPartitionBytes参数、使用coalesce()和repartition()方法,但性能提升甚微,怀疑是Spark的「小文件问题」导致,寻求在AWS Glue环境下大幅缩短处理时间的优化方案。

当前实现代码:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("ProcessJSON").getOrCreate()

# Reading JSON files directly with Spark
df = spark.read.json(f"s3://{bucket}/{prefix}*.json")

# Processing and transformations
df.show()
优化方案

1. 用Spark wholeTextFiles读取并合并小文件

Spark直接读取JSON时,每个小文件会生成单独分区,导致Task数量过多(10万+),调度和执行开销极大。wholeTextFiles可将多个小文件打包到一个RDD分区,大幅减少Task数量后再统一解析JSON:

from pyspark.sql import SparkSession
import json

spark = SparkSession.builder.appName("ProcessJSON").getOrCreate()

# 根据总数据量设置目标分区数(5GB数据建议20-50个)
target_partitions = 30
rdd = spark.sparkContext.wholeTextFiles(f"s3://{bucket}/{prefix}*.json", minPartitions=target_partitions)

# 解析单条JSON记录
df = rdd.map(lambda file_content: json.loads(file_content[1])).toDF()

# 后续处理
# df.show()  # 生产环境建议删除该调试操作

2. 使用AWS Glue DynamicFrame优化小文件读取

Glue的DynamicFrame针对S3小文件场景做了原生优化,通过groupFiles参数可自动合并小文件,减少分区数:

import sys
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions

args = getResolvedOptions(sys.argv, ['JOB_NAME'])
glueContext = GlueContext(spark.sparkContext)
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# 读取时启用小文件合并,设置合并后单文件大小为100MB
dynamic_frame = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={
        "paths": [f"s3://{bucket}/{prefix}"],
        "recurse": True,
        "groupFiles": "inPartition",
        "groupSize": "104857600"  # 100MB
    },
    format="json"
)

# 转换为DataFrame继续处理
df = dynamic_frame.toDF()
# df.show()

job.commit()

3. 提前合并S3小文件(预处理)

如果允许预处理,可先将S3上的小JSON文件合并为大文件(比如100MB/个),彻底解决小文件问题:

方法1:AWS CLI批量合并(示例框架)

# 遍历前缀下的文件,每1000个合并为一个大文件
aws s3 ls s3://your-bucket/path/ --recursive | awk '{print $4}' > file_list.txt
split -l 1000 file_list.txt batch_

for batch in batch_*; do
    target_file="merged_$(date +%Y%m%d_%H%M%S).json"
    # 循环读取batch中的文件并合并写入目标路径
    while read file; do
        aws s3 cp s3://your-bucket/$file - >> temp.json
    done < $batch
    aws s3 cp temp.json s3://your-bucket/merged-path/$target_file
    rm temp.json
done

方法2:AWS Lambda自动合并

编写Lambda函数遍历指定S3前缀,将小文件合并为大文件,适合自动化预处理场景。

4. 调整Glue作业核心配置参数

  • 升级到Glue 4.0(Spark 3.3):新版本针对Spark文件读取、任务调度做了大量优化,相比3.0可提升20%-30%处理速度。
  • 优化Spark分区参数:在Glue作业的「作业参数」中添加以下配置:
    • --conf spark.sql.shuffle.partitions=30:根据总数据量设置,5GB数据建议20-50
    • --conf spark.hadoop.mapreduce.input.fileinputformat.split.minsize=104857600:强制Spark合并小文件为100MB的分区
    • --conf spark.hadoop.mapreduce.input.fileinputformat.split.maxsize=104857600
  • 选择高性能DPU:使用G.2X或G.4X DPU,这类实例拥有更多CPU和内存,更适合IO密集型的小文件处理任务。

5. 移除冗余操作

  • 删除df.show():该操作会触发全量数据计算,仅用于调试,正式运行需删除或替换为df.count()等轻量校验操作。
  • 减少中间数据落地:直接在DataFrame上执行后续写入或转换操作,避免不必要的磁盘IO开销。

内容的提问来源于stack exchange,提问作者Jayron Soares

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 20:34:51