Apache Spark 1.6.2在YARN集群的资源与时长限制方案咨询
嘿,针对你在YARN集群上用Spark 1.6.2遇到的问题——那些简单查询霸占所有资源还动不动崩溃的情况,你提到的三个限制需求都有可行的方案,我来逐一拆解:
1. 限制任务运行时长
你可以从SQL查询级别和整个应用级别分别设置超时:
- SQL查询级别的超时控制:Spark 1.6支持
spark.sql.queryTimeout参数(单位:秒),专门用来限制单个SQL查询的运行时间。比如提交任务时加上:
这样任何运行超过1小时的SQL查询都会被自动中断,完美解决--conf spark.sql.queryTimeout=3600SELECT * FROM 1TB表这种无意义的长时间查询。 - YARN应用级别的超时:如果要限制整个Spark应用的生命周期,不管是SQL还是RDD任务,都可以用YARN的应用超时机制。提交任务时添加:
这里--conf spark.yarn.appMasterEnv.YARN_APP_TIMEOUT=3600000 --conf spark.yarn.maxAppAttempts=1YARN_APP_TIMEOUT是毫秒级(示例为1小时),搭配spark.yarn.maxAppAttempts=1可以避免超时后任务自动重试,进一步减少资源浪费。也可以在YARN全局配置yarn-site.xml中设置yarn.app.timeout,对所有应用生效。
2. 限制Shuffle磁盘空间
Shuffle磁盘占用过高的问题,可以结合Spark自身优化和YARN的磁盘配额来解决:
- Spark层面的Shuffle优化:
- 调整内存溢出阈值:
spark.shuffle.spill.memoryThreshold(默认5MB),当Shuffle数据在内存中超过这个值时会刷到磁盘,适当调小可以避免内存过载,同时控制磁盘写入的节奏; - 开启压缩:确保
spark.shuffle.spill.compress=true和spark.shuffle.compress=true(默认已开启),用压缩算法减少Shuffle文件的体积; - 优化缓冲区:调整
spark.shuffle.file.buffer(默认32KB),增大缓冲区可以减少小文件数量,间接降低磁盘占用。
- 调整内存溢出阈值:
- YARN容器磁盘配额限制:YARN可以强制限制每个容器的磁盘使用量。在
yarn-site.xml中设置yarn.nodemanager.container-disk-mb,或者提交Spark任务时指定:
这里的--conf spark.yarn.executor.memoryOverhead=2048memoryOverhead包含了executor的磁盘、内存等额外资源,当容器磁盘使用超过配额时,YARN NodeManager会直接终止该容器,从根源上限制Shuffle的磁盘占用。另外,建议把spark.local.dir指向多个磁盘目录,分散Shuffle的磁盘负载。
3. 限制单查询/用户CPU资源(类似Circuit Breaker)
这个需求可以从资源隔离和主动中断两个方向入手,实现类似Elasticsearch的circuit breaker效果:
- YARN队列级别的资源隔离:
- 如果你用Capacity Scheduler,在
capacity-scheduler.xml里配置:yarn.scheduler.capacity.<queue-path>.capacity:队列的最小资源占比,保证每个队列有基础资源;yarn.scheduler.capacity.<queue-path>.maximum-capacity:队列的最大资源占比,防止单个队列抢占所有资源;yarn.scheduler.capacity.<queue-path>.user-limit-factor:单个用户在队列中能使用的资源倍数,比如设为1,意味着单个用户最多占用该队列的全部资源,避免某一个用户独占集群。
- 如果你用Fair Scheduler,在
fair-scheduler.xml中给每个用户或队列设置maxResources,限制其能使用的最大CPU核数和内存,比如:<user name="dev"> <maxResources>8 cores, 16384mb</maxResources> </user>
- 如果你用Capacity Scheduler,在
- Spark应用级别的CPU限制:提交任务时通过
spark.cores.max限制整个应用的总CPU核数,比如:
这样不管应用怎么调度,最多只能用10个CPU核,不会霸占整个集群。同时用--conf spark.cores.max=10spark.executor.cores限制每个executor的核数,控制并行度。 - 主动中断(自定义Circuit Breaker):Spark 1.6没有内置的circuit breaker,但可以通过自定义
SparkListener实现。你可以写一个监听类,在任务启动时记录时间,定时检查查询的运行时长或CPU累计使用时间,当超过阈值时调用SparkContext.cancelAllJobs()终止任务。比如监听onJobStart和onJobEnd事件,跟踪每个用户的任务运行时间,达到上限就主动中断。
内容的提问来源于stack exchange,提问作者Thomas Decaux
相关产品推荐
相关产品推荐

