如何在DSE中配置Spark Job Server实现作业排队等待资源?
DSE Spark作业排队与并行数控制方案
针对你在DSE Cassandra环境下,需要限制Spark作业并行数为4、超出作业自动排队的需求,可通过以下几种方式实现:
一、调整Spark Job Server的并行作业配置
Spark Job Server本身支持配置单个Spark上下文的并发作业数,修改后超出的作业会自动进入内部队列等待,而非直接拒绝:
- 找到Spark Job Server的配置文件(通常为
job-server.conf或集成在dse.yaml中),添加/修改以下参数:- 若针对预创建的特定Spark上下文:
spark.jobserver.contexts.your-precreated-context.concurrent-jobs = 4 - 若设置全局默认的上下文并发数:
spark.jobserver.default.context.concurrent-jobs = 4
- 若针对预创建的特定Spark上下文:
- 同时确保全局作业数限制不低于此值,可设置:
spark.jobserver.max-jobs-per-context = 4 - 修改完成后重启Spark Job Server,新提交的作业会在并行数达到4时进入队列,待已有作业完成后自动调度执行。
二、Spark集群资源参数适配
结合你的3节点OLAP集群配置(8核/32GB节点),需调整Spark资源参数,确保4个并行作业不会超出集群承载:
- 针对预创建的Spark上下文,在创建时指定以下参数:
# 集群总可用CPU核数(3节点×8核) spark.cores.max = 24 # 每个executor分配的核数,确保4个作业可均分资源 spark.executor.cores = 6 # 每个executor分配的内存(预留8GB给Cassandra与系统) spark.executor.memory = 24G # 调度模式设为FIFO,确保队列作业按提交顺序执行 spark.scheduler.mode = FIFO - 这些参数可在创建预创建上下文时通过API传入,或配置在Spark Job Server的上下文定义中。
三、客户端自定义排队逻辑(进阶方案)
若Spark Job Server内置队列无法满足更灵活的调度需求,可在作业提交客户端实现本地排队逻辑:
- 核心思路:维护一个作业队列,实时监控运行中的作业数量,仅当并行数低于4时提交新作业,否则将作业加入队列等待。
- 示例伪代码(Python):
import requests import time from queue import Queue from threading import Thread JOB_SERVER_API = "http://your-job-server-address:8090" MAX_PARALLEL = 4 job_queue = Queue() running_job_ids = set() def check_running_jobs(): """定期检查运行中作业的状态,清理已完成的作业ID""" while True: completed = [] for job_id in running_job_ids: resp = requests.get(f"{JOB_SERVER_API}/jobs/{job_id}") status = resp.json().get("status") if status in ["FINISHED", "FAILED", "KILLED"]: completed.append(job_id) for job_id in completed: running_job_ids.remove(job_id) time.sleep(5) def submit_job(job_payload): """提交单个作业,确保不超出并行数限制""" while len(running_job_ids) >= MAX_PARALLEL: time.sleep(3) resp = requests.post(f"{JOB_SERVER_API}/jobs", json=job_payload) job_id = resp.json().get("jobId") running_job_ids.add(job_id) return job_id def process_queue(): """处理队列中的作业""" while True: if not job_queue.empty() and len(running_job_ids) < MAX_PARALLEL: payload = job_queue.get() submit_job(payload) time.sleep(3) # 启动监控与队列处理线程 Thread(target=check_running_jobs, daemon=True).start() Thread(target=process_queue, daemon=True).start() # 提交作业时调用此方法 def add_job_to_queue(payload): job_queue.put(payload)
四、DSE资源管理器补充配置
在dse.yaml中调整Spark资源分配,确保集群为Spark预留足够资源:
spark: executor_instances: 3 executor_cores: 8 executor_memory: "24G"
此配置为每个节点分配1个executor,每个executor使用8核与24GB内存,总资源刚好支撑4个并行作业的均分需求。
内容的提问来源于stack exchange,提问作者Prateek Sharma
相关产品推荐
相关产品推荐

