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

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
已尝试方案
  1. 推测S3前缀目录结构导致倾斜(部分前缀仅含1个文件,部分含数千个),于是用哈希码扁平化目录结构,使每个文件夹仅存一个文件:

原目录结构:

/parent1
    /file1
    /file2
    ...
    /file1000

/parent2/
    /file1

修改后结构:

hashcode=FEFRE#$#$$#FE/parent1/file1
hashcode=#$#$#Cdvfvf@#/parent1/file1

但此操作无效果。

  1. 使用超大集群,试图靠足够内存抵消输入倾斜,问题仍存在。

观察发现每个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 22:30:48