如何优化处理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
相关产品推荐
相关产品推荐

