news 2026/8/3 11:09:11

Spark 核心之 Driver 原理剖析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Spark 核心之 Driver 原理剖析

摘要:如果把 Spark 应用比作一个人,Driver 就是它的大脑。从 spark-submit 敲下回车的那一刻起,SparkContext 初始化、DAGScheduler 切分 Stage、TaskScheduler 分发 Task、SchedulerBackend 与集群通信、SparkEnv 管理运行时环境——所有这些都在 Driver 内部精密协作。本文从 Driver 内部架构全景、SparkContext 初始化链路、三大调度器协作模型、SparkEnv 七大组件、Driver 生命周期六个维度,配合 1 张原创深色架构图和完整源码级分析,带你彻底看清 Driver 的内部世界。

关键词:Spark Driver, SparkContext, DAGScheduler, TaskScheduler, SchedulerBackend, SparkEnv, BlockManager, Driver 生命周期


一、开篇:Driver 是什么?

先回答一个面试高频题:

“Spark Driver 到底做了什么?”

答案不是一句话能说完的。Driver 是 Spark 应用的总控制器,它承载了至少以下七大职责:

#职责核心组件
1解析用户代码 → 构建 DAGDAGScheduler
2将 DAG 切分为 StageDAGScheduler
3将 Stage 拆分为 Task 并分发TaskSchedulerImpl
4与集群通信(申请/释放资源)SchedulerBackend
5管理运行时环境(内存/序列化/Shuffle)SparkEnv
6事件监听与 Web UILiveListenerBus + SparkUI
7Executor 心跳监控HeartbeatReceiver

二、Driver 内部架构全景图

2.1 三大组件群

┌─────────────────────────────────────────────────┐ │ SparkContext │ │ ┌─────────────┐ ┌─────────────┐ ┌──────────┐│ │ │核心调度组件 │ │ SparkEnv │ │ 监控/事件 ││ │ │DAGScheduler │ │BlockManager │ │LiveListen ││ │ │TaskScheduler│ │ShuffleMgr │ │SparkUI ││ │ │SchedulerBknd│ │MemoryMgr │ │MetricsSys ││ │ └─────────────┘ └─────────────┘ └──────────┘│ └─────────────────────────────────────────────────┘

三、SparkContext 初始化链路

这是 Driver 启动最核心的代码路径。

// 源码:SparkContext.scala (简化版初始化链路)classSparkContext(config:SparkConf)extendsLogging{// Step 1: 创建 SparkEnv(运行时环境)privatevar_env:SparkEnv=_ _env=SparkEnv.createDriverEnv(conf,isLocal,listenerBus,...)// Step 2: 创建元数据追踪器_applicationId=_env.conf.get("spark.app.id")_dagScheduler=newDAGScheduler(this)// Step 3: 创建 TaskScheduler + SchedulerBackendval(sched,ts)=SparkContext.createTaskScheduler(this,master,deployMode)_schedulerBackend=sched _taskScheduler=ts// Step 4: DAGScheduler 绑定 TaskScheduler_dagScheduler=newDAGScheduler(this)_taskScheduler.start()// Step 5: 启动心跳接收器_heartbeatReceiver=env.rpcEnv.setupEndpoint(HeartbeatReceiver.ENDPOINT_NAME,newHeartbeatReceiver(this))// Step 6: 注册 SparkListener + 启动 WebUIsetupAndStartListenerBus()_ui=SparkUI.create(conf,listenerBus,_env,...)}

四、三大调度器协作模型 🔥

这是 Driver 最核心的调度链路。

用户代码 (Action) │ ▼ DAGScheduler.handleJobSubmitted() │ ① 回溯 RDD 依赖 → 创建 ResultStage │ ② getMissingParentStages() → 递归构建 ShuffleMapStage │ ③ submitMissingTasks() → 为每个 Partition 创建 Task ▼ TaskScheduler.submitTasks(taskSet) │ ④ TaskSetManager 封装 → 数据本地性排序 │ ⑤ reviveOffers() → 通知 Backend 有 Task 可调度 ▼ SchedulerBackend.reviveOffers() │ ⑥ makeOffers() → 匹配空闲 Executor 与 Task │ ⑦ launchTasks() → 序列化 Task 发送给 Executor ▼ Executor (远程 JVM) ⑧ 反序列化 → 执行 → 序列化结果 → StatusUpdate

4.1 DAGScheduler:Stage 切分核心

// 源码核心逻辑privatedefsubmitStage(stage:Stage):Unit={valmissing=getMissingParentStages(stage).sortBy(_.id)if(missing.isEmpty){submitMissingTasks(stage,jobId.get)}else{for(parent<-missing)submitStage(parent)}}privatedefgetMissingParentStages(stage:Stage):List[Stage]={stage.rdd.dependencies.flatMap{caseshufDep:ShuffleDependency[_,_,_]=>getOrCreateShuffleMapStage(shufDep,stage.firstJobId)case_=>Nil// NarrowDep 不切分}.toList}

规则:遇到 ShuffleDependency 即切分 Stage。

4.2 TaskSchedulerImpl:数据本地性

// 数据本地性优先级PROCESS_LOCAL>NODE_LOCAL>RACK_LOCAL>ANY// 每个级别等待 spark.locality.wait (默认 3s)

4.3 SchedulerBackend:集群通信适配器

实现通信目标
StandaloneSchedulerBackendSpark Master (Netty RPC)
YarnSchedulerBackendYARN AM → RM (Hadoop RPC)
KubernetesClusterSchedulerBackendK8s API Server (HTTP REST)

五、SparkEnv:运行时环境七大组件

// 源码:SparkEnv.createDriverEnv()valblockManager=newBlockManager(...)valbroadcastManager=newBroadcastManager(...)valmapOutputTracker=newMapOutputTrackerMaster(...)valshuffleManager=SortShuffleManager(conf)valmemoryManager=UnifiedMemoryManager(conf,...)valserializer=newJavaSerializer(conf)// or KryoSerializervalclosureSerializer=newJavaSerializer(conf)
组件职责
BlockManagerRDD 缓存(Memory + Disk)、Shuffle 数据存储
MapOutputTracker追踪 Shuffle Map 输出位置(Master/Worker)
ShuffleManagerSortShuffleManager 管理 Shuffle 写/读
MemoryManagerUnifiedMemoryManager:执行 + 存储统一内存池
SerializerTask 序列化/反序列化
BroadcastManagerTorrentBroadcast 分布式广播
RpcEnvNettyRpcEnv:Driver ↔ Executor 通信基础设施

六、Driver 完整生命周期

Phase 1: spark-submit → main() → new SparkContext() ├── 创建 SparkEnv(运行时环境) ├── 创建 DAGScheduler + TaskScheduler + SchedulerBackend ├── 向 Master/RM 注册,申请 Executor └── 启动 HeartbeatReceiver + SparkUI Phase 2: Action 触发 → DAG 调度 ├── DAGScheduler.handleJobSubmitted() ├── Stage 切分 + Task 生成 ├── TaskScheduler 分发 Task └── Executor 执行 + StatusUpdate 回传 Phase 3: 监控与运维 ├── LiveListenerBus 推送事件 ├── SparkUI :4040 实时监控 └── HeartbeatReceiver 心跳检测 Phase 4: sc.stop() → 优雅退出 ├── 通知 SchedulerBackend 停止 ├── Kill 全部 Executor ├── 向 Master/RM 注销 └── 释放 SparkEnv 资源

七、Driver 配置调优

spark-submit\--driver-memory 4G\# Driver JVM 堆内存--driver-cores2\# Driver 可用核心--confspark.driver.maxResultSize=2G\# collect() 结果上限--confspark.driver.extraJavaOptions="-XX:+UseG1GC"\--confspark.driver.extraClassPath=/path/to/extra.jar\--confspark.driver.supervise=true\# Standalone Cluster 专属my-app.jar

八、总结

要点总结
Driver 本质用户 main() 运行的 JVM 进程,SparkContext 即 Driver 入口
三大调度器DAGScheduler → TaskScheduler → SchedulerBackend 逐层下发
SparkEnv7 大组件提供序列化、Shuffle、内存、存储等运行时能力
生命周期初始化 → 调度循环 → 监控 → 优雅退出

金句:Executor 是 Spark 的四肢,Driver 是 Spark 的大脑——DAGScheduler 思考"如何拆分",TaskScheduler 决定"派给谁",SchedulerBackend 负责"怎么送"。


作者:starzy | AI Data Engineer / 大数据技术实践者
博客:blog.starzy.cn | GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/3 11:09:05

扬声器内磁盖模具设计要点与精密成型技术

1. 扬声器内磁盖模具设计概述扬声器内磁盖是扬声器磁路系统中的关键部件&#xff0c;其作用是固定磁体并形成均匀的磁场间隙。模具设计质量直接影响内磁盖的尺寸精度、机械强度和电磁性能。在实际项目中&#xff0c;一套优秀的内磁盖模具需要同时满足精密成型、高效生产和成本控…

作者头像 李华
网站建设 2026/8/3 11:05:48

Keyviz完全指南:免费开源键鼠可视化工具让你的操作一目了然

Keyviz完全指南&#xff1a;免费开源键鼠可视化工具让你的操作一目了然 【免费下载链接】keyviz Keyviz is a free and open-source tool to visualize your keystrokes ⌨️ and &#x1f5b1;️ mouse actions in real-time. 项目地址: https://gitcode.com/gh_mirrors/ke/…

作者头像 李华
网站建设 2026/8/3 11:05:08

信息系统架构设计与优化实战指南

1. 信息系统基础认知第一次接触信息系统这个概念是在2008年&#xff0c;当时我参与了一个零售企业的库存管理系统升级项目。那套老旧的系统每天要处理上万条出入库记录&#xff0c;经常在业务高峰期崩溃。项目经理指着那台嗡嗡作响的服务器说&#xff1a;"这就是我们企业的…

作者头像 李华
网站建设 2026/8/3 11:05:03

SpringBoot教师业绩管理系统开发实践

1. 项目概述 教师业绩管理系统是高校信息化建设中的重要组成部分&#xff0c;主要用于记录、统计和分析教师的教学、科研等工作成果。基于SpringBoot框架开发的教师业绩管理系统&#xff0c;能够有效解决传统手工记录方式效率低下、数据易丢失、统计困难等问题。 我在实际开发…

作者头像 李华
网站建设 2026/8/3 11:04:53

Scroll Reverser:5分钟解决macOS滚动方向冲突的终极方案

Scroll Reverser&#xff1a;5分钟解决macOS滚动方向冲突的终极方案 【免费下载链接】Scroll-Reverser Per-device scrolling prefs on macOS. 项目地址: https://gitcode.com/gh_mirrors/sc/Scroll-Reverser 你是否经常在Mac触控板和鼠标之间切换&#xff0c;却因为滚动…

作者头像 李华