TT Lab
开始
学习 学习路径 课程

Apache Spark — 慢作业的答案在执行计划和事件日志里

转换只累积计划,行动才启动作业并把它拆分为阶段和任务

在 TT Lab 中继续学习

一句话总结

Spark 代码看上去像是自上而下执行的,但实际上是驱动器一边堆叠计划,一遇到行动算子(action),就把整个计划转换成作业(job),再拆成阶段和任务分给执行器去运行。什么事在什么时候发生,不是由代码、而是由事件日志来告诉你。

为什么这样划分

第一次读 Spark 代码时,很容易以为 spark.read.csv(...) 这一行在读文件,filter 这一行在过滤,count 这一行在计数。这样读的话,就会在错误的地方寻找性能问题。看上去“读取花了 40 秒”的那个位置,其实是前面堆叠起来的二十个转换算子一起运行的位置。

Spark 之所以这样划分工作,是因为它从数据装不进一台机器的内存这个前提出发。要把计算分给多个进程,就得有人了解全局、把工作切碎分发,而拿到碎片的一方只需要计算自己的那一份。集群模式概述中的术语表是这样定义这些角色的。

作业的定义里包含“为响应行动算子”这句话,这是关键。没有行动算子,就没有作业。

工作原理

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 文件,后一个阶段才能读取它,所以阶段在这里分开。

用 groupBy 计数的一个作业被分成两个阶段的示意图。第一个阶段中,每个输入分区的一个任务连续完成 CSV 读取、过滤、部分聚合,并写出 shuffle 文件。第二个阶段的任务从所有 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 模式下碰巧看得到,但搬到集群上就会消失。

实际工作中真正重要的事

下一项实验要做什么

用 spark-submit 确认版本,然后启动数一数 30 万条订单的第一个作业。运行一个只堆叠了转换算子的应用,确认事件日志里一个作业也没有;再在一个多次调用了行动算子的应用里,用日志数一数实际启动了几个作业。用 groupBy 引发 shuffle,观察阶段被拆开;并比较在 local[1] 和 local[2] 下,同一个文件的分区数有什么不同。最后读取一个不存在的路径,拿到错误条件名,并把这些数字整理成报告。