Spark 核心之 Master 原理剖析
摘要Spark Master 是 Standalone 模式的集群大脑。它管理所有 Worker 的注册与心跳调度 Driver 和 Application 的资源分配通过 SpreadOutApps 算法确保负载均衡借助 ZooKeeper 实现主备选举和状态恢复。本文从 Master 内部架构全景、核心数据结构、schedule() 调度引擎源码、Worker 生命周期管理、HA 高可用机制五个维度配合 1 张原创深色架构图 完整源码分析带你彻底看清 Master 的内部世界。关键词Spark Master, schedule(), Worker 管理, 资源调度, SpreadOutApps, HA, ZooKeeper, LeaderElection, PersistenceEngine一、开篇Master 是什么如果说 Driver 是单个应用的大脑那么 Master 就是整个 Standalone 集群的指挥中心。职责说明Worker 管理注册、心跳、超时检测、资源登记资源调度Driver 调度 Application Executor 资源分配应用生命周期注册 → 排队 → 调度 → 完成 → 移除HA 高可用ZK 主备选举、状态持久化、故障自动切换一句话定位Master 不执行 Task不做 DAG 解析——它只做一件事谁需要资源、谁有资源、把资源分给谁。二、Master 内部架构全景图Master 四大核心模块┌────────────────────────────────────────────┐ │ Spark Master │ │ ┌──────────────┐ ┌──────────────────┐ │ │ │ 核心数据结构 │ │ schedule() 引擎 │ │ │ │ workers: Set │ │ SpreadOutApps │ │ │ │ apps: Set │ │ FIFO/FAIR 策略 │ │ │ │ drivers: Set │ │ launchExecutor() │ │ │ └──────────────┘ └──────────────────┘ │ │ ┌──────────────┐ ┌──────────────────┐ │ │ │ RPC Endpoint │ │ HA 高可用 │ │ │ │ RegisterWorker│ │ LeaderElection │ │ │ │ RegisterApp │ │ PersistenceEngine│ │ │ │ Heartbeat │ │ 状态恢复 │ │ │ └──────────────┘ └──────────────────┘ │ └────────────────────────────────────────────┘三、核心数据结构// 源码Master.scalaprivate[master]classMaster(overridevalrpcEnv:RpcEnv,address:RpcAddress,webUiPort:Int,valsecurityMgr:SecurityManager,valconf:SparkConf)extendsThreadSafeRpcEndpoint{// 核心数据结构 valworkersnewHashSet[WorkerInfo]// 已注册的 WorkervalappsnewHashSet[ApplicationInfo]// 运行中的应用valdriversnewHashSet[DriverInfo]// 已提交的 Driver (Cluster 模式)// 等待队列 valwaitingAppsnewArrayBuffer[ApplicationInfo]// 等待资源分配的应用valwaitingDriversnewArrayBuffer[DriverInfo]// 等待 Worker 的 Driver// 已完成的应用数量限制用于 Web UI 展示privatevalcompletedAppsnewArrayBuffer[ApplicationInfo]privatevalcompletedDriversnewArrayBuffer[DriverInfo]}数据结构类型用途workersHashSet所有已注册的 WorkerappsHashSet运行中的 ApplicationdriversHashSetCluster 模式提交的 DriverwaitingAppsArrayBuffer排队等待资源分配的应用waitingDriversArrayBuffer排队等待 Worker 的 Driver四、schedule() 调度引擎 这是 Master 最核心的方法——每次状态变更都会触发调用。// 源码Master.scala - schedule() 核心逻辑privatedefschedule():Unit{// Step 1: 先调度 Driver (Cluster 模式)for(driver-waitingDrivers.toList){valworkerworkers.filter(_.stateWorkerState.ALIVE).filter(canLaunchDriver(_,driver.desc)).headOption worker.foreach{wwaitingDrivers-driver w.endpoint.send(LaunchDriver(driver.id,driver.desc))driver.stateDriverState.RUNNING}}// Step 2: 再调度 Application (Executor 分配)startExecutorsOnWorkers()}4.1 SpreadOutApps 负载均衡算法// 源码startExecutorsOnWorkers()privatedefstartExecutorsOnWorkers():Unit{for(app-waitingAppsifapp.coresLeft0){valusableWorkersworkers.filter(_.stateWorkerState.ALIVE).filter(canLaunchExecutor(_,app.desc))// SpreadOut: 尽可能将 Executor 分散到不同 Worker// 每个 Worker 分配尽可能少的 cores最少 1 个 ExecutorvalnumWorkersusableWorkers.lengthvarcoresPerWorkerapp.desc.maxCores.getOrElse(Int.MaxValue)/numWorkersif(coresPerWorker1)coresPerWorker1for(worker-usableWorkersifapp.coresLeft0){valcoresToUsemath.min(worker.coresFree,coresPerWorker)launchExecutor(worker,app,coresToUse)}}}SpreadOut 策略的优势将任务分散到更多节点 → 更高的数据本地性命中率 → 减少 Shuffle 网络开销。五、Worker 生命周期管理Worker 启动 → RegisterWorker → Master 登记 → Heartbeat (每 15s) ↓ 超时未心跳? / \ NO YES 继续运行 RemoveWorker ↓ 状态标记 DEAD waitingApps 重新排队 apps 中的 Executor 失联// Worker 心跳超时检测caseHeartbeat(workerId,worker)worker.lastHeartbeatSystem.currentTimeMillis()// 15 秒超时检查在 checkForWorkerTimeOuts() 中级联效应Worker 宕机 →RemoveWorker→ 释放的资源重新进入schedule()→ 其他 Worker 承接任务。六、Master HAZooKeeper 主备选举# 启动 HA Master 集群# Master 1./sbin/start-master.sh-hmaster-1 --webui-port8080\--confspark.deploy.recoveryModeZOOKEEPER\--confspark.deploy.zookeeper.urlzk1:2181,zk2:2181,zk3:2181# Master 2 (备)./sbin/start-master.sh-hmaster-2 --webui-port8080\--confspark.deploy.recoveryModeZOOKEEPER\--confspark.deploy.zookeeper.urlzk1:2181,zk2:2181,zk3:21816.1 选举机制// 源码ZooKeeperLeaderElectionAgent.scalaclassZooKeeperLeaderElectionAgent(valmasterInstance:LeaderElectable,conf:SparkConf,securityMgr:SecurityManager,zkUrl:String)extendsLeaderElectionAgent{// 在 ZK 创建临时顺序节点 /spark/leader_election/lock-0000000xxx// 序号最小的节点成为 Leader// Leader 宕机 → 临时节点自动删除 → 下一个最小序号成为新 Leader}6.2 状态持久化与恢复// 恢复模式sealedtraitRecoveryModecaseobjectZOOKEEPERextendsRecoveryModecaseobjectFILESYSTEMextendsRecoveryMode// 本地文件系统caseobjectCUSTOMextendsRecoveryMode// 自定义实现caseobjectNONEextendsRecoveryMode// 无 HA恢复流程新 Leader 从持久化存储读取 → 恢复 workers/apps/drivers 状态 → 执行schedule()继续调度。七、总结要点总结定位Master 集群资源调度中心不执行 Task/DAG核心方法schedule()Driver Executor 资源分配调度策略SpreadOutApps 分散负载 FIFO/FAIR 排队数据结构workers/apps/drivers/waitingApps/waitingDriversHAZK 临时节点选举 持久化状态恢复金句如果 Driver 是单个应用的大脑Master 就是整个集群的交管中心——它不跑车但它决定哪辆车走哪条路。作者starzy | AI Data Engineer / 大数据技术实践者博客blog.starzy.cn | GitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践