Spark读取S3非可拆分嵌套JSON遇任务倾斜与内存问题求助
问题背景与现象
处理存储在S3中的高度嵌套、不可拆分JSON数据集,文件大小差异极大——最小10KB,最大可达300MB。执行以下代码读取文件并repartition到目标分区数时,出现任务倾斜:多数任务数秒内完成,但有一个任务耗时数小时后触发内存问题(心跳丢失/堆空间不足等)。尝试repartition是为了打乱分区与文件的映射关系,因为Spark可能按顺序读取文件,且同一目录下的文件特性趋同(要么全大,要么全小)。
df = spark.read.json('s3://my/parent/directory') df.repartition(396)
相关配置:
- default parallelism = 396
- 总核心数 = 400
已尝试方案
- 推测S3前缀目录结构导致倾斜(部分前缀仅含1个文件,部分含数千个),于是用哈希码扁平化目录结构,使每个文件夹仅存一个文件:
原目录结构:
/parent1 /file1 /file2 ... /file1000 /parent2/ /file1
修改后结构:
hashcode=FEFRE#$#$$#FE/parent1/file1 hashcode=#$#$#Cdvfvf@#/parent1/file1
但此操作无效果。
- 使用超大集群,试图靠足够内存抵消输入倾斜,问题仍存在。
观察发现每个Spark分区分配的文件数在2到32之间(因JSON不可拆分,每个文件对应DataFrame一行),怀疑Spark依据spark.sql.files.maxPartitionBytes分配文件到分区——大文件所在分区仅分配2个文件,小文件所在分区分配更多文件?
补充:最新尝试及异常
按照建议调整后的完整代码如下,作业卡在Running状态,Worker日志显示任务已完成,但Driver日志报错“can not increase buffer size”(触发2GB缓冲区限制问题),未执行collect等会拉取数据到Driver的操作。
本次作业配置(确保Executor每个任务有足够内存):
- Driver/Executor核心数:2
- Driver/Executor内存:198 GB
# 示例S3源路径: s3://masked-bucket-urcx-qaxd-dccf/gp/hash=8e71405b-d3c8-473f-b3ec-64fa6515e135/group=34/db_name=db01/2022-12-15_db01.json.gz ##################### # 初始化代码 ##################### spark = SparkSession.builder.appName('query-analysis').getOrCreate() ##################### # 查询Dataframe ##################### dfq = spark.read.json(S3_STL_RAW, primitivesAsString="true") dfq.write.partitionBy('group').mode('overwrite').parquet(S3_DEST_PATH)
需求
在无法修改输入文件大小的前提下,求能让作业正常运行、任务均匀分布的解决方案。
内容的提问来源于stack exchange,提问作者Deepak Gaur
相关产品推荐
相关产品推荐

