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

PySpark读取大表时Job未生成、Executor异常及SQL查询延迟问题求助

问题分析与解决方案

从你遇到的现象和贴出来的代码来看,这个问题核心是Spark在生成Job前的元数据扫描、分片计算阶段卡了太久,还顺带搞挂了大部分Executor,下面是具体的分析和解决建议:

一、先搞懂为什么会出现这些现象

你说SQL标签有个运行的Query,但Jobs标签看不到Job,这是因为Spark还没到提交Job的阶段——Query停留在「解析元数据、生成数据分片」的准备环节,这个过程耗了近1小时才启动Job。而大部分Executor因为长时间等不到任务、资源过载或者和Driver通信出问题,被标记成了死亡状态,只剩2个撑着。

二、针对性的解决办法

1. 优化分区表的分片策略,减少元数据扫描耗时

你查的是分区表table_name,虽然指定了pt范围,但Spark得扫描这些分区里的所有ORC文件元数据来生成分片,这一步最容易卡:

  • 换个ORC分片策略:你当前用的ETL策略会给每个小ORC文件单独生成一个分片,如果你的表每个分区有大量小文件,分片数量会爆炸,元数据处理直接卡死。改成BI策略,它会合并小文件的分片,能大幅减少分片数:
    .config('hive.exec.orc.split.strategy','BI')
    
  • 降低分片并行度:你设的spark.sql.getSplit.parallelism.num是180,这个参数控制元数据扫描的并行线程数,过高会把Driver和Executor的资源占满,导致Executor挂掉。先降到60,和你配置的Executor实例数匹配试试:
    .config('spark.sql.getSplit.parallelism.num', '60')
    

2. 检查YARN队列的资源够不够养60个Executor

你配置了60个Executor,每个要12G内存+3核CPU,但你用的DataBusinessDept_Risk_Query队列可能根本没这么多资源:

  • 找集群管理员确认队列的总配额:如果队列总内存远小于60*12G=720G,总CPU远小于60*3=180核,YARN要么启动不了这么多Executor,要么启动后直接把超标的给kill了,最后只剩2个能活下来。
  • 调整Executor配置:要是队列资源不够,先把spark.executor.instances降到30左右,同时可以把单Executor内存提到16G,用「少实例大内存」的方式适配队列资源,反而更稳定。

3. 调整缓存的姿势,别让它拖后腿

你现在的代码是先查SQL再cache,然后用count触发计算,其实可以换个顺序,先让Spark完成元数据处理和分片,再缓存:

# 先触发计算,让Spark走完元数据和分片阶段
a = spark.sql(""" select t2.* from table_name t2 where pt between '20210401' and '20210416' """)
a.count()
# 再缓存,此时数据已经分片完成,缓存更高效
a.cache()

或者直接用SQL的缓存提示,让Spark更智能地处理:

spark.sql(""" CACHE TABLE temp_table AS select t2.* from table_name t2 where pt between '20210401' and '20210416' """)
a = spark.table("temp_table")

4. 查日志搞清楚Executor到底为啥死了

要彻底解决Executor死亡的问题,得看日志:

  • 去YARN的ResourceManager UI找到你的应用,点开每个Executor的日志,看看是不是内存不够(OutOfMemoryError)、和Driver通信断了,或者被YARN强制kill了。
  • 把Driver的日志级别调高点,在代码里加一行:
    spark.sparkContext.setLogLevel("WARN")
    
    看看Driver日志里有没有关于Executor断开连接的报错,能帮你精准定位。

三、额外的小优化

你代码里加了动态分区的配置,但你当前的代码只是读表,根本没写分区表,这些配置可以先删掉,减少Spark初始化的开销:

# 这两行暂时注释掉,等需要写分区表的时候再加
# .config("spark.hadoop.hive.exec.dynamic.partition", "true").\
# .config("spark.hadoop.hive.exec.dynamic.partition.mode", "nonstrict").\

内容的提问来源于stack exchange,提问作者muzhen xv

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 16:17:34