有朋友跑来问我:他们公司刚把Kyuubi和Spark 3.4.1这套东西搭起来,用同一个账号从beeline往Kyuubi里提SQL,跑了一段时间发现YARN队列怎么都用不满——队列明明显示还有大把空闲资源,可任务就是跑得慢悠悠的,executor数量也一直上不去。这个问题非常典型,尤其是把Kyuubi作为统一SQL网关、所有任务都以单用户身份提交的场景里。今天我把排查思路和最终结论完整写出来,从架构、动态资源分配、参数作用域、YARN调度器到SQL并行度,一层一层剥开,希望对正在被同样问题折磨的人有帮助。
1. 先看清楚Kyuubi单用户模式下,YARN队列里的资源到底是谁在占用
1.1 "单用户提交任务"背后的资源结构
很多人对Kyuubi的资源模型是模糊的。Kyuubi本身只是一个网关服务,它不直接跑计算。你通过beeline或者JDBC连上Kyuubi,它会在YARN上拉起一个Spark engine,这个engine才是一个真正的Spark Application。换句话说,你在YARN ResourceManager页面上看到的那个application,并不是你提交的那条SQL,而是Kyuubi为这个用户启动的长驻Spark引擎。
默认情况下Kyuubi的引擎共享级别是USER,也就是同一个用户的所有session会复用同一个engine。单用户场景下,无论你开几个beeline连接,背后基本都是同一个Spark Application、同一个SparkContext。SQL提进去之后,SparkContext在已有的executor上调度task,不会再为每条SQL去YARN申请独立的资源。
所以"YARN队列使用不满"这个现象,其实对应的是:这个常驻engine占用队列资源的情况没有达到你的预期。这里有个容易被忽略的点——你看到的不是"任务"的用量,而是"引擎"的用量。这条认知是整个排查的地基。
1.2 "使用不满"的三种典型外观
实际看YARN队列的时候,"不满"至少有三种完全不同的表象,对应完全不同的排查方向:
| 外观 | 现象 | 优先排查方向 |
|---|---|---|
| 容器数量少 | 队列空闲资源很多,但application的executor数量一直没涨上去 | Spark动态资源分配、SQL并行度 |
| 容器规格不对 | executor是有的,但每个executor的vcore/内存比预期小,总量上不去 | Spark executor参数、YARN单容器上限 |
| 容器数够但task跑不满 | executor数量不少,但stage里的task只有几十个,大量executor空转 | AQE合并、shuffle partition、数据倾斜 |
我看到太多人一上来就改YARN队列配额,其实很多"不满"根本不是队列给的资源不够,而是Spark端压根没申请、或者申请了没跑满。先分清这三种情况,后面才不会白忙活。
登录ResourceManager页面,找到对应application,点开看"Application"概览里的容器数量和资源量。如果发现Pending的容器请求很多,那是"申请了但没拿到";如果Pending是0,那是"根本不想多要",问题在Spark这一侧。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 动态资源分配:单用户场景下的"第一疑犯"
2.1 动态分配是怎么决定何时向YARN多要container的
Spark on YARN里,如果不开动态资源分配,executor数量就是固定的,你配了几个就是几个,多了浪费少了不够。开了动态分配之后,Spark的ExecutorAllocationManager会周期性地检查任务积压情况,然后决定加不加executor。
加executor的触发条件简单说就是:当前有pending task(task已经生成但还没有executor去跑),同时集群还有资源余量。管理线程每隔spark.dynamicAllocation.schedulerBacklogTimeout秒看一次,发现有积压就追加executor,直到满足需求或者达到spark.dynamicAllocation.maxExecutors上限。反过来,如果executor空闲超过spark.dynamicAllocation.executorIdleTimeout,就会被回收。
这就引出一个特别关键的反直觉结论:如果任务本身生成不出足够多的task,动态分配不但不会加executor,甚至会觉得"当前规模已经够用"。 我把这个逻辑打了个比方:动态分配不是"队列里还有资源就抢过来用",而是"发现有人排队才开新窗口"。没排队,窗口自然不会加。所以你看到一个空荡荡的队列,第一反应不应该是去调YARN,而是去数一下任务到底有多少个task。
2.2 Kyuubi引擎复用机制让动态分配更容易踩坑
Kyuubi场景下动态分配还有它特殊的一面。因为engine是常驻的、复用的,两条SQL之间的空窗期里,executor会被动态分配回收掉。等到下一条SQL进来,又要从1个、2个executor重新扩展。这个"扩容滞后"在短查询、交互式查询上尤其明显:SQL跑完了executor还没加到理想数量,自然看起来队列一直是半满状态。
另外Spark 3.x在YARN上开动态分配,基本都会遇到shuffle文件归属的问题。当一个executor被动态回收时,如果它磁盘上还有shuffle中间文件,依赖这些文件的后续task就会出问题。旧方案是部署ExternalShuffleService,让executor即使被删除shuffle文件也在;新方案更省事,直接开spark.dynamicAllocation.shuffleTracking.enabled=true,让Spark跟踪每个executor上的shuffle文件,等文件过期后再回收executor。
在Kyuubi环境里,我基本都建议开shuffleTracking而不是额外部署ESS。但这里有个坑:spark.dynamicAllocation.shuffleTracking.timeout如果设置太大,executor会被shuffle文件"绑架"很久,明明任务跑完了一直不释放,队列看起来又是"用不满但也不释放"的状态。取一个合理的超时,比如300到600秒,平衡文件保留和资源回收。
2.3 先照着一张清单排查动态分配参数
下面这份参数清单是我在处理"Kyuubi+Spark on YARN队列用不满"时必查的,建议直接对照:
| 参数 | 默认值 | 建议值 | 说明 |
|---|---|---|---|
| spark.dynamicAllocation.enabled | false | true | 不开这个,executor数量固定,基本告别弹性 |
| spark.dynamicAllocation.shuffleTracking.enabled | false | true | 避免部署ESS,同时保证shuffle不丢 |
| spark.dynamicAllocation.initialExecutors | minExecutors | 1 | 引擎刚启动时先占住的最小数量 |
| spark.dynamicAllocation.minExecutors | 0 | 1或2 | 太小则每轮SQL都要从零扩容 |
| spark.dynamicAllocation.maxExecutors | 无限 | 按队列配额算 | 见第4节的公式 |
| spark.dynamicAllocation.schedulerBacklogTimeout | 1s | 5s | 扩得太激进会导致抖动,太慢则任务等资源 |
| spark.dynamicAllocation.sustainedSchedulerBacklogTimeout | 1s | 30s | 连续积压后下一次扩容的间隔 |
| spark.dynamicAllocation.executorIdleTimeout | 60s | 120s | 太短会造成反复扩缩的震荡 |
| spark.dynamicAllocation.shuffleTracking.timeout | Long.Max | 300s~600s | 太长则executor长期被shuffle文件占住 |
先别急着照抄,逐个确认一遍当前值。排查的时候用spark.dynamicAllocation.enabled=false就能解释大多数"队列一直不满且executor数量恒定偏小"的现象。
3. 逐层缩小问题范围:从YARN到SparkContext再到SQL的完整排查链路
3.1 第一层:判断是"没申请"还是"没分配"
站点里来了问题,我习惯从最外层往里面推,第一站永远是ResourceManager页面。
找到目标application,看两个数据:一个是当前已分配的容器,另一个是pending的资源请求。如果pending的容器数量大于0,说明Spark已经在向YARN要资源了,是YARN这边给不出来。这时候要看队列的容量配置、单容器上限、用户限制,可能是队列配了小气点。
如果pending是0,说明SparkContext当前压根不认为需要更多executor,问题就直接越过YARN,落到Spark调度和SQL这边。从经验看,八成"单用户+Kyuubi队列不满"最后都落在这条分支。
3.2 第二层:Spark UI的Executors页签会告诉你真相
接下来打开Spark UI,切到Executors页签。看两个东西:Total Executors和Active Executors。动态分配打开时,Total会随着负载变化,Active则表示当前真正在跑task的executor。如果两者差距大,说明一批executor被拉起来了但没事干,任务并行度跟不上。
再往上看每个executor的"Task Time"和"Active Tasks"。如果所有executor的Active Tasks都是0,而stage又没结束,大概率是pending task太少,跑完了之后executor在傻等。如果某些executor的Task Time明显高于其他,那是数据倾斜的典型信号,后面第5节单独讲。
还有一个入口容易被忽略:Executors页签里动态分配的"Remove"行为。回收executor时,这里会出现被移除的executor记录,结合executorIdleTimeout参数能判断是不是回收策略太激进,导致引擎下一轮SQL又要重新扩容。
3.3 第三层:用SQL计划反推"任务就只需要这么多资源"
很多"队列不满"其实是SQL并行度不够造的,跟前两层都没关系。判断方法很简单:在Kyuubi里执行EXPLAIN,或者直接在Spark UI的SQL页签里看每个stage的task数量。
举个例子,一次SELECT COUNT(*) FROM big_table,如果表是128MB一个文件块,5GB的表大约40个split,scan stage最多也就生成几十个task。这个stage跑完之后,count汇总阶段只剩一个task。这种情况下你就算队列里能塞下50个executor,实际用的也就几个,剩下全部空转。这是SQL本身决定的,不是资源分配的bug。
别跟资源调度较劲,先问SQL到底要不要这么大并行度。判断方法非常直接:如果某个stage的max task数小于你期望的executor总核数,那"队列用不满"就是这个SQL的正常行为,想让它跑得更满,得去提高task数量,而不是去调YARN。
4. 参数作用域与配置传递:为什么你改了参数,Kyuubi没生效
4.1 Kyuubi里的三种配置注入路径
很多人在Kyuubi上改参数,改完发现没变化,其实是没有搞清楚Kyuubi的配置作用域。Kyuubi的配置有三种注入路径:
- 引擎级配置:写在
kyuubi-defaults.conf里的spark.*配置,engine启动的时候一次性生效,影响整个SparkContext。 - 连接级配置:JDBC URL后面跟的
;spark.xxx=yyy参数,在engine首次创建时生效,如果engine已经被别的连接创建了,这些参数不会重新作用。 - 会话级配置:连接之后执行
SET spark.xxx=yyy,改的是当前session的运行时配置,部分Spark参数是支持动态更新的。
这是一个非常常见的坑:你改了JDBC URL的参数,但Kyuubi里的engine已经因为之前的连接存在了,新参数根本没有被加载,你还以为配置生效了。排查时建议用一个新的engine来验证,或者临时把共享级别调一下。
4.2 你需要合理计算单用户场景的资源上限
看队列"不满"之前,先算一下理论上这个单用户最多能吃多少资源。这个上限由三个因素共同决定,取最小值:
- YARN队列的最大可用容量,由队列的
maximum-capacity决定; - 容量调度器里单用户限制,默认
user-limit-factor=1,表示单个用户最多使用到队列guaranteed capacity对应的资源量; - Spark动态分配上限,由
spark.dynamicAllocation.maxExecutors * spark.executor.cores决定。
假设整集群400核,YARN队列capacity配了30%,maximum-capacity是60%。那么队列最多能拿到240核,但用户因子默认是1,单用户实际上限最保守只有120核——也就是capacity的30%这个数。如果你把maxExecutors * executor.cores配到200,YARN也不会真给你200,多出来的容器要么pending要么干脆不分配。这也能解释很多"我明明配了50个executor,怎么只起来了20个"的疑问。
还有一个经常被忽略的量:YARN的单容器上限。yarn.scheduler.maximum-allocation-mb和yarn.scheduler.maximum-allocation-vcores决定了每个container能给多大。你如果把spark.executor.cores配成8,而YARN单容器vcore上限是4,那container请求就永远不会被满足。从RM页面上看,它永远处于pending状态,队列也是"怎么都用不满"。
4.3 一个实际配置案例(可直接抄)
假设你确认了队列单用户实际能给到120核,我建议这样配:
code复制spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.shuffleTracking.enabled=true
spark.dynamicAllocation.initialExecutors=1
spark.dynamicAllocation.minExecutors=1
spark.dynamicAllocation.maxExecutors=28
spark.dynamicAllocation.schedulerBacklogTimeout=5s
spark.dynamicAllocation.executorIdleTimeout=120s
spark.dynamicAllocation.shuffleTracking.timeout=300s
spark.executor.cores=4
spark.executor.memory=8g
spark.executor.memoryOverhead=2g
spark.cores.max=112
算一下:单个executor占用YARN container资源为8g+2g=10g内存、4个vcore。28个executor就是112核,加上driver本身占一些,刚好不超过120核上限。注意driver和AM也会占1个container,所以executor数量不能直接顶满上限。spark.cores.max设成112是为了跟executor配置对齐,防止某些场景下Spark按总核数算超额申请。
5. 顺着任务特征继续挖:AQE、Shuffle Partition和数据倾斜
5.1 Spark 3.4.1 AQE默认开启,它对资源使用的影响
Spark 3.4.1里spark.sql.adaptive.enabled默认是true,AQE会在运行时根据实际shuffle数据量动态调整reduce端的分区数。它默认开着coalescePartitions,也就是把小分区合并,减少task数量。
这个优化在大多数时候是好事,但它会让"队列不满"的现象更明显。比如一个大join,本来shuffle产生1000个分区,AQE一看很多分区数据量很小,直接把分区合并成60个,于是后续stage只有60个task。配合60个task,你就算有30个executor、120核,也只有60个slot在干活。这时候单纯把spark.dynamicAllocation.maxExecutors调大是没用的,因为task数不够,动态分配压根不会扩张。
遇到这种情况,不要急着骂AQE。先看执行计划里有没有出现CustomShuffleReader coalesced这样的字样,如果有,说明AQE确实在做分区合并。这里给两个方向:如果这批SQL是短平快的交互查询,保持AQE合并收益更大;如果是批量ETL、希望吃满集群资源,可以试试spark.sql.adaptive.coalescePartitions.parallelismFirst=true,让AQE优先保留更高并行度,或者直接把spark.sql.adaptive.coalescePartitions.enabled=false关掉。
5.2 shuffle partition 设置对"队列使用不满"的传导
AQE之外,spark.sql.shuffle.partitions也是一个直接变量。默认值是200,意思是所有shuffle类算子(group by、join、distinct)的reduce阶段最多生成200个task。
假设你按上一节的方案配了28个executor、4核,总共112核,理论上能同时跑112个task。但shuffle partitions是200,那么reduce阶段要跑两波:第一波112个task,第二波88个task,后面88个task跑的时候有24个核空转。如果你数据量大、shuffle段耗时很长,两波之间的资源浪费就很明显。这时候把这个值调到500、800甚至更高,让每个task的数据量保持在合理范围,同时又尽量吃满核数。
但这里永远有个平衡:shuffle partition调得太大,会产生太多小文件,影响后续读取效率,也可能让driver端的调度压力变大。我的建议是:先跑一次看AQE实际合并到了多少分区,再以数据量除以期望的单个task处理量(单task处理50-100MB是经验范围)来估算,不要拍脑袋乱填。
5.3 数据倾斜造成"局部忙碌,整体不满"
还有一种情况是,"不满"只是表象,实际上是某些executor在扛大梁。比如一个key的数据量占了80%,所有数据都shuffle到同一个executor上,这个executor跑得要死,其他executor早就跑完在空等。从YARN队列总览看,资源占用率不高,因为绝大多数executor处于空闲待命状态,但从任务进度看,却又卡着不动。
这种问题在Spark UI的Stage页面非常明显:executor的task时间柱状图长短差异巨大,某个executor的task时间远超中位数。如果是AQE下的join倾斜,可以打开spark.sql.adaptive.skewJoin.enabled(3.4.1默认开启的,但需要确认生效);如果是group by倾斜,AQE帮不上忙,需要用加盐、两阶段聚合这类手段自己处理。
个人体会是,数据倾斜和资源不满经常同时出现,很多人的第一反应是加资源,加完之后倾斜的task还是那么慢,资源更浪费了。先搞清楚是"资源不够"还是"数据不均衡",再动手调。
6. 验证与监控:优化之后如何证明队列真的吃满了
6.1 用YARN RM页面和Spark UI联合判读
改完参数别急着下结论,先建立一套判断标准。我一般看三个地方:
第一是RM页面上application的已分配容器和pending请求。如果优化后pending变成0、已分配数明显提高,说明Spark端确实在要资源,而且YARN给了。
第二是Spark UI的Executors页签。观察一个完整SQL周期里,executor总数是否先快速爬升、跑完后逐步回落。如果还是趴在地平线上不动,说明动态分配或者并行度还有问题。
第三是SQL页签里每个stage的task量。这个最能说明"为什么不满"。如果每个stage的task数都远小于总核数,那不是资源问题,是SQL并行度问题;如果task数足够,但executor没上来,才轮到动态分配背锅。
6.2 一个简单的对比实验
想验证猜测到底对不对,我习惯做对比实验。比如怀疑是shuffle partitions过低导致不满,锁定一条固定的聚合SQL,分别用200和800跑一遍,对比Spark UI里reduce阶段task数量和executor峰值数。注意Kyuubi的engine是复用缓存的,改参数之前要确保是新建的engine,否则新参数不生效。
同一张表、同一条SQL,把spark.sql.shuffle.partitions从200改到800之后,如果executor峰值从10个爬到了28个,说明之前的根因就是task数太少。如果改完executor还是10个,那就要回到动态分配参数和YARN配额上继续排查。用这样的对比法,能快速把问题钉死在某一层。
还要提醒一句:做这种实验最好挑非高峰时段,并且临时关闭或者隔离其他任务,避免队列里别的应用干扰判断。队列资源是动态的,其他app占多占少会直接影响你的实验结论。
6.3 我踩过的两个坑
最后分享两个实战中踩过的坑,都是这次"Kyuubi+Spark 3.4.1队列不满"排查里真实遇到的。
第一个坑是盲目调大executor规格。最初我把spark.executor.memory从8g调成16g、cores保持4,结果executor一直起不来,RM页面里一堆pending container。原因是YARN单容器内存上限是12g,16g内存加上overhead早就超了。后来把executor内存调回8g、overhead设成2g,容器立刻正常铺满。资源申请最怕的就是"规格超出YARN上限",一定是集群配置倒推executor规格,而不是反过来。
第二个坑是shuffleTracking超时设置太激进。有段时间我把spark.dynamicAllocation.shuffleTracking.timeout设成了60秒,想着加快回收,结果SQL跑一半executor就被删了,后面stage重新拉shuffle数据,任务反而更慢。后来调成300秒,给shuffle文件留够生命周期,动态分配才稳定下来。这个参数本质是让"历史包袱"和"资源释放"做平衡,太小容易误删,太大就会看到executor占着不做事。
这类问题的排查,我现在的习惯是:先看RM,再开Spark UI,最后用EXPLAIN复盘SQL;先确认"该不该满",再谈"怎么让它满"。很多时候,让一个SQL跑满队列并不是目标,让队列资源确实用在了它需要的并行度上,跑得快、不浪费,才是真正值得调的。
