Apache Spark — 慢作业的答案在执行计划和事件日志里
转换只累积计划,行动才启动作业并把它拆分为阶段和任务
一句话总结
Spark 代码看上去像是自上而下执行的,但实际上是驱动器一边堆叠计划,一遇到行动算子(action),就把整个计划转换成作业(job),再拆成阶段和任务分给执行器去运行。什么事在什么时候发生,不是由代码、而是由事件日志来告诉你。
为什么这样划分
第一次读 Spark 代码时,很容易以为 spark.read.csv(...) 这一行在读文件,filter 这一行在过滤,count 这一行在计数。这样读的话,就会在错误的地方寻找性能问题。看上去“读取花了 40 秒”的那个位置,其实是前面堆叠起来的二十个转换算子一起运行的位置。
Spark 之所以这样划分工作,是因为它从数据装不进一台机器的内存这个前提出发。要把计算分给多个进程,就得有人了解全局、把工作切碎分发,而拿到碎片的一方只需要计算自己的那一份。集群模式概述中的术语表是这样定义这些角色的。
- 驱动器(driver):运行应用的 main() 并创建 SparkContext 的进程。负责制定计划并分发任务。
- 执行器(executor):为每个应用在 Worker 节点上启动的进程,负责运行任务,并把数据保存在内存或磁盘中。
- 任务(task):发送给一个执行器的工作单位。
- 作业(job):为响应一个行动算子而产生的、由多个任务组成的并行计算。
- 阶段(stage):把作业拆成相互依赖的、更小的任务集合。文档写道,它与 MapReduce 的 map、reduce 阶段类似。
作业的定义里包含“为响应行动算子”这句话,这是关键。没有行动算子,就没有作业。
工作原理
RDD 编程指南明确写道“Spark 中的所有转换都是惰性的”。转换算子不会立即计算结果,只记住施加了哪些转换,等到行动算子要求把结果返回给驱动器时才会计算。文档也写出了这种设计的好处:如果预先知道 map 之后会有 reduce,就可以只把缩减后的结果送回驱动器,而不是庞大的中间结果。看到终点之后再选路,这就是惰性的价值。
df = spark.read.csv("/data/shop/orders.csv", header=True, schema=ddl) # 계획 1줄
paid = df.filter(F.col("status") == "paid") # 계획 +1
big = paid.withColumn("bulk", F.col("qty") >= 5) # 계획 +1
big.count() # 행동 — 여기서 잡이 뜨고, 위 세 줄이 한 번에 실행된다
遇到行动算子时,驱动器会把计划转换成物理计划,并在 shuffle 边界处把它切开。没有 shuffle 的区段(读取、过滤、添加列这样的窄转换)可以在一个任务内连续运行,所以是一个阶段。当出现像 groupBy 这样需要把相同的键汇集到一处的宽转换时,前一个阶段的所有任务必须先把结果写成 shuffle 文件,后一个阶段才能读取它,所以阶段在这里分开。
在阶段内部,一个分区就是一个任务。所以分区数就是可以同时工作的份数。读取文件时,分区数并不只由文件大小决定。根据性能调优文档,切分文件时的最小分区数(spark.sql.files.minPartitionNum)的默认值是 spark.sql.leafNodeDefaultParallelism,而该值的默认值是 SparkContext 的默认并行度。这意味着,即使是小文件,也会尝试至少按核心数来切分读取。
在 local 模式下,什么是驱动器,什么是执行器
这个实验 Pod 里没有集群。在 Master URL 表中,local 是一个工作线程(没有并行),local[K] 是 K 个工作线程。在 local 模式下,一个驱动器 JVM 同时承担执行器的角色。任务在那个 JVM 内部的线程里运行。
即便如此,驱动器、执行器、作业、阶段、任务的区分依然存在。只不过制定计划的代码和运行任务的代码在同一个进程里而已。这就是用 local 模式学习并不白费的原因。事件日志中留下的作业、阶段、任务事件的形态,和在集群中是一样的。这个 Pod 被设置成 local[2],驱动器内存为 1g。
事件日志——回溯已结束应用的证据
驱动器的 Web UI(4040)只有在应用存活期间才能看到。要查看已结束的应用,按监控文档的说明使用历史服务器,而历史服务器是用应用留下的事件日志重建 UI。没有事件日志,对已结束的应用几乎什么都无从得知。
事件日志是每行一个 JSON 对应一个事件。作业开始时会记录 SparkListenerJobStart,阶段结束时会记录 SparkListenerStageCompleted。所以像“三次行动就有三个作业”这样的猜测,可以通过数行数来验证。
有一点需要注意。根据配置文档,Spark 4.2 的默认值是:spark.eventLog.compress 为 true,压缩编解码器为 zstd,spark.eventLog.rolling.enabled 为 true。保持默认的话,日志会被滚动成多个文件并被压缩,无法直接用 grep 读取。本实验已把它改成了单个纯文本文件。如果想在生产环境里用同样的方法,要先确认这两项配置。
在现场相遇的样子
一个行动算子并不等于一个作业。如果不提供 schema 就读取带表头行的 CSV,Spark 为了确认表头行,会先启动一个小作业。写文件的行动算子也可能产生不止一个作业。所以不要数代码,要数日志。本实验中,对于调用了三个行动算子的应用,也不是猜测作业数,而是直接从日志里数。
出错的位置有两种。路径错了,不会等到行动算子,在 read 的那一刻就会在分析阶段失败。这时的错误条件名是 PATH_NOT_FOUND。相反,如果问题出在数据里的值,在制定计划时是无从得知的,要到行动算子真正读取数据时才会爆出来。前者的代码行和错误是对得上的,后者则是在无关的那一行(行动算子)上爆出来。
驱动器会先死掉。为了用眼睛看结果而调用 collect(),全部结果会汇集到驱动器这一处。RDD 指南指出,这可能让驱动器内存溢出,所以只看几条的时候建议用 take。无论执行器有多少个,驱动器只有一个。
转换算子里的 print 不会出现在驱动器的屏幕上。因为转换算子里的代码是在执行器上运行的。在 local 模式下碰巧看得到,但搬到集群上就会消失。
实际工作中真正重要的事
- 转换算子是计划,行动算子才让它干活。慢的地方通常是带有行动算子的那一行,但原因是它前面的全部转换算子。
- 阶段的边界就是 shuffle。减少阶段数,就是减少 shuffle。
- 一个分区就是一个任务。分区数决定了可以同时工作的上限。
- 作业数不要猜,要从事件日志里数。schema 推断、表头行确认、写入都会多出作业。
- 事件日志的默认值是滚动和 zstd 压缩。选择读取工具之前,先看配置。
- 驱动器只有一个。用 collect 把结果汇集到驱动器,是最常见的驱动器内存事故。
下一项实验要做什么
用 spark-submit 确认版本,然后启动数一数 30 万条订单的第一个作业。运行一个只堆叠了转换算子的应用,确认事件日志里一个作业也没有;再在一个多次调用了行动算子的应用里,用日志数一数实际启动了几个作业。用 groupBy 引发 shuffle,观察阶段被拆开;并比较在 local[1] 和 local[2] 下,同一个文件的分区数有什么不同。最后读取一个不存在的路径,拿到错误条件名,并把这些数字整理成报告。