EMR上PySpark缓存DataFrame执行count()极慢问题排查
问题根因
该问题是配置错误、Spark机制理解偏差、数据倾斜三类问题叠加导致,核心原因按影响权重排序:
1. Executor GC参数存在致命笔误,未启用预期的G1收集器
你配置的spark.executor.extraJavaOptions中,GC启用参数写为-X:+seG1GC,属于无效配置:
- 正确的G1启用参数为
-XX:+UseG1GC,当前配置不仅参数前缀错误(-X:应为-XX:),参数名也拼写错误,JVM启动时会直接忽略该参数,Executor实际运行在默认的Parallel GC上 - 你配置的
XX:InitiatingHeapOccupancyPercent=35是G1专属参数,对Parallel GC不生效。18G堆内存下,Parallel GC触发Full GC时单轮停顿可达数分钟到数十分钟,GC线程会占满CPU核心,业务代码执行几乎停滞,这就是你观测到单核心100%占用、其余核心空闲的直接原因
2. 全局存储级别配置无效,缓存内存占用远超预期
你设置的spark.storage.level: MEMORY_AND_DISK_SER不是Spark的有效配置项,Spark不存在全局指定存储级别的参数:
- 存储级别必须在代码中调用
persist()方法时传入,直接调用cache()默认使用非序列化的MEMORY_AND_DISK级别 - 非序列化缓存的内存占用是序列化缓存的3~5倍,直接抬升Executor内存压力,进一步拉长GC停顿时间
3. 动态分配+数据本地性调度导致任务单点排队
你开启了动态资源分配且设置spark.executor.instances=0,触发了Spark的本地性调度限制:
- Cache作业执行阶段,动态分配会申请足够的Executor资源并行完成缓存写入,因此10分钟的缓存构建耗时符合预期
- 缓存完成后,空闲Executor会被动态分配机制回收,但缓存块会持久化在原节点的本地磁盘/残留内存中
- 执行
count()时,Spark调度优先遵循NODE_LOCAL本地性原则,只会将对应分区的任务调度到持有缓存块的节点上。你观测到的5个卡住的任务,对应的缓存块全部落在同一个节点上,而该节点仅分配到1个可用Executor核心,5个任务只能单线程排队执行 - 默认数据本地性等待时间较长,Spark不会轻易将任务降级调度到无缓存块的节点,导致集群其余核心全程空闲
4. 分区倾斜放大故障影响
原始Parquet数据均匀、Join阶段运行正常,不代表Join后分区均匀:空值会被Shuffle拉取到同一个分区,最终导致5个分区的数据量是普通分区的数十到上百倍,这也是这几个任务运行速度远低于其余195个秒级完成任务的原因。
修复方案
- 修正GC参数:将Executor Java选项中的错误参数
-X:+seG1GC替换为-XX:+UseG1GC,和Driver侧GC配置保持一致 - 修正缓存逻辑:删除无效的全局存储级别配置,代码中将
cache()替换为persist(StorageLevel.MEMORY_AND_DISK_SER),缓存调用后立即执行一次count()触发缓存全量物化,避免后续作业触发缓存计算 - 优化动态分配与本地性配置:
- 设置
spark.dynamicAllocation.minExecutors=12(约为集群总核数的20%),避免缓存完成后所有Executor被回收 - 将
spark.locality.wait调整为3s,本地性等待超时后自动将任务调度到其他可用节点,避免单点排队
- 设置
- 解决分区倾斜:缓存前对DataFrame调用
repartition(256)(分区数取集群总vcore数的1~2倍),打散倾斜的大分区,保证所有分区数据量均匀 - 对齐并行度配置:设置
spark.sql.shuffle.partitions=256,和全局并行度匹配,避免Join阶段产生过少分区导致倾斜
内容的提问来源于stack exchange,提问作者Guillaume
相关产品推荐
相关产品推荐

