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的日志级别调高点,在代码里加一行:
看看Driver日志里有没有关于Executor断开连接的报错,能帮你精准定位。spark.sparkContext.setLogLevel("WARN")
三、额外的小优化
你代码里加了动态分区的配置,但你当前的代码只是读表,根本没写分区表,这些配置可以先删掉,减少Spark初始化的开销:
# 这两行暂时注释掉,等需要写分区表的时候再加 # .config("spark.hadoop.hive.exec.dynamic.partition", "true").\ # .config("spark.hadoop.hive.exec.dynamic.partition.mode", "nonstrict").\
内容的提问来源于stack exchange,提问作者muzhen xv
相关产品推荐
相关产品推荐

