运行在Amazon EMR上的Spark应用Executor输入量为何大于实际处理文件大小
核心原因
你遇到的Executor页输入指标大于实际文件大小、SQL页指标和文件大小匹配的问题,本质是两个页面对「输入数据量」的统计口径完全不同,前者统计异常偏高直接对应作业慢的根因:
- SQL标签页的输入指标,统计的是Spark SQL解析层最终拿到的有效业务数据大小,也就是源文件解压、过滤掉无效读取内容后的实际处理数据量,因此和你感知的源文件真实大小完全对齐。
- Executor监控页的输入指标,统计的是Executor进程从存储层(S3/HDFS/本地盘)实际读取的所有字节总量,包含重复读取、预读冗余、传输开销等所有非有效业务数据的流量,正常场景下这个值会比有效数据大10%~30%,如果超出这个范围甚至翻倍,结合你当前的集群配置,基本是以下几个问题导致:
- 资源严重超卖引发频繁重试。2台c5.24xlarge单台规格为96vCPU、192GiB内存,扣除系统和EMR基础服务占用后,YARN可调度资源约为180vCPU、360GiB。你同时跑20个Spark应用,每个分配3个executor+1个driver,总共要调度80个进程,就算给每个executor/driver只分配2vCPU、4GiB内存,总资源需求也达到160vCPU、320GiB,几乎占满所有可用资源,没有给OS页缓存、Shuffle、Python进程(如果有)留冗余。这种场景下Executor会频繁因为OOM、内存超用被YARN杀掉,Task失败重试会重复拉取同一份源数据,这部分重复流量全部会计入Executor输入指标,但SQL层只会统计最终成功Task读取的有效数据,不会统计重试的冗余流量。
- 不可切分压缩文件导致冗余读取。如果你的源文件用GZIP、ZIP这类不支持并行切分的压缩格式存储,Spark为了实现并行计算,会让多个Task重复拉取同一个压缩块的完整内容,再在本地解压后过滤出自己需要处理的分片,这部分重复拉取的流量只会计入Executor侧的统计,不会算入SQL层的有效输入。
- 对象存储预读/重试带来的冗余流量。如果数据存放在S3上,EMR默认的S3A连接器会开启顺序预读机制,每个Task读取数据时会提前拉取后续块的内容到本地缓存,如果预读的内容没被当前Task实际处理,这部分流量会被算入Executor输入;如果遇到S3连接超时、读取中断触发客户端重试,重复拉取的流量也会被统计到Executor输入中。
- 数据本地性失效带来的额外传输开销。2台Worker上混部60个Executor,YARN调度时很容易把Task分配到没有存储对应数据块的节点,跨节点拉取数据块时的网络传输副本、校验开销都会被算入Executor的输入流量,不会计入SQL层的有效数据统计。
排查与解决步骤
按优先级从高到低排查:
- 先解决资源超卖问题
- 登录YARN ResourceManager UI查看集群整体vCPU、内存使用率,再到Spark作业的Event Log里统计每个Stage的Task失败重试率、Executor非正常退出次数,如果单Stage重试率超过5%、存在Executor被OOM Kill的记录,先收缩资源配额:单台c5.24xlarge建议给每个Spark应用最多分配2个executor,每个executor配置4vCPU、8GiB内存,driver配置2vCPU、4GiB内存,总资源占用控制在集群可调度资源的70%以内,留出足够的冗余给系统缓存、Shuffle和后台进程。如果20个应用需要同时跑,建议扩容Worker节点,不要硬超资源。
- 检查源文件格式与压缩配置
- 确认源文件的压缩格式,如果使用GZIP等不可切分压缩算法,替换为Snappy、ZSTD、LZO这类支持并行切分的压缩算法,同时把单文件大小调整到和Spark分区大小匹配(建议128MiB~256MiB),避免单文件过大导致的重复读取。
- 检查
spark.sql.files.maxPartitionBytes参数,不要将该值设置得远小于存储层块大小(S3/HDFS默认块大小128MiB),否则会产生过多小分区,引发大量冗余的文件打开、读取开销。
- 优化对象存储读取配置(S3场景适用)
- 关闭S3A的默认预读机制:设置
spark.hadoop.fs.s3a.prefetch.enabled=false,或者将spark.hadoop.fs.s3a.readahead.range调整为1MiB以内,减少无意义的预读冗余流量。 - 针对分析类作业的随机读场景,设置
spark.hadoop.fs.s3a.experimental.input.fadvise=random,关闭顺序读优化逻辑,减少不必要的预读。 - 适当调大S3客户端的超时时间,避免因为网络波动触发客户端重连重读。
- 关闭S3A的默认预读机制:设置
- 优化数据本地性
- 查看Spark UI中各Stage的Task本地性分布,如果超过30%的Task本地性级别为RACK_LOCAL或ANY,调大
spark.locality.wait参数到5s,让调度器等待Task分配到持有对应数据块的节点后再启动,减少跨节点拉取的流量开销。 - 如果使用HDFS存储数据,将块副本数设置为2,保证两个Worker节点都持有对应块的副本,提升本地调度命中率。
- 查看Spark UI中各Stage的Task本地性分布,如果超过30%的Task本地性级别为RACK_LOCAL或ANY,调大
内容的提问来源于stack exchange,提问作者sriparth
相关产品推荐
相关产品推荐

