Apache Spark — 慢作业的答案在执行计划和事件日志里
启动第一个作业并用事件日志确认——转换是惰性的,行动才触发工作
目标
用 local 模式的 Spark 读取并计数 30 万条订单,并通过事件日志比较只堆叠了转换算子的应用和调用了行动算子的应用。用数字确认作业、阶段、任务何时产生,以及分区数由什么决定。
为什么重要
第一次读 Spark 代码时,每一行看上去都像是在原地立即执行。其实不是。filter、withColumn、select 这样的转换算子只是在计划上多加一行,到了调用 count、take、write 这样的行动算子时,整个计划才被转换成作业并执行。不了解这一点,就会忽略一件事:看上去“读取很慢”的地方,其实是前面堆叠起来的全部转换算子一起运行的地方。
这个差别可以用眼睛确认。Spark 在应用结束时,会把它做过的事情以事件日志的形式留下来。启动了几个作业、有几个阶段、有几个任务、用了什么物理计划,全都以每行一个 JSON 的形式记录在里面。在生产环境里回溯已结束的作业时看的也是它(历史服务器会读取这个文件)。
在 local 模式下,一个驱动器 JVM 就是执行器。local[2] 的意思是使用两个核心,这个数字同时决定了默认并行度和文件被切成几块。
步骤
- 把
spark-submit --version的输出(包括标准错误)保存到 /root/spk/first/version.txt。 - 创建 /root/spk/first/count.py,以应用名
spk-first-count统计/data/shop/orders.csv(带表头行)的行数,并以{"rows": 정수}(占位符为整数)写入 /root/spk/first/out/count.json。 - 创建 /root/spk/first/lazy.py,以应用名
spk-first-lazy自己提供 schema 进行读取,堆叠filter(或where)和withColumn,但不要调用行动算子。运行之后,该应用的作业必须为 0 个。 - 创建 /root/spk/first/actions.py,以应用名
spk-first-actions调用三次以上行动算子(例如count、take、write),并把 100 行以 JSON 写入 /root/spk/first/out/sample。在该应用的事件日志里数一数作业开始事件的个数,以一个整数写入 /root/spk/first/out/jobs.txt。 - 创建 /root/spk/first/agg.py,以应用名
spk-first-agg把按status统计的订单数以带表头行的 CSV 写入 /root/spk/first/out/by_status。把该应用中已完成的阶段数以整数写入 /root/spk/first/out/stages.txt。 - 创建 /root/spk/first/par.py,让它把第一个参数作为应用名,用
--master local[1]运行spk-first-par1,用--master local[2]运行spk-first-par2,分别把master、default_parallelism、partitions(读取订单文件得到的 DataFrame 的分区数)写入 /root/spk/first/out/spk-first-par1.json 和 /root/spk/first/out/spk-first-par2.json。 - 创建 /root/spk/first/fail.py,以应用名
spk-first-fail读取不存在的路径/data/shop/order.csv,把捕获到的异常的错误条件名(getCondition())写到 /root/spk/first/out/error.txt 的第一行。 - 在 /root/spk/first/report.md 中以
## 잡과 스테이지(韩文,意为“作业与阶段”)、## 게으른 실행(韩文,意为“惰性执行”)、## 파티션(韩文,意为“分区”)三节写下你的数字。第一节放入第 4 步的作业数,第三节放入第 6 步的两个分区数。
参考
- 运行方式像
spark-submit /root/spk/first/count.py这样。默认配置在/opt/spark/conf/spark-defaults.conf中(local[2]、驱动器内存 1g、事件日志已开启)。 - 事件日志在应用结束后会留在
/root/spark-events/local-<숫자>(占位符为数字)。运行过程中,名字末尾带有.inprogress。哪个文件对应哪个应用,可以用grep -l '"App Name":"spk-first-actions"' /root/spark-events/*来找。 - 作业数用
grep -c '"Event":"SparkListenerJobStart"' <파일>(占位符为文件)来数,阶段数则数"Event":"SparkListenerStageCompleted"。也可以用jq -c 'select(.Event=="SparkListenerJobStart")'查看内容。 - 常见错误:漏掉
spark.stop(),导致日志以.inprogress的状态留下;对带表头行的 CSV 漏掉header=True,把表头行当成了一行数据来数;不提供 schema 就读取,产生了确认表头行的作业(第 3 步必须是 0 个作业)。 - 官方文档:Cluster Mode Overview · RDD Programming Guide — RDD Operations · Monitoring — Viewing After the Fact · Submitting Applications — Master URLs
确认是哪个 Spark
把 spark-submit --version 的输出连同标准错误一起保存到 /root/spk/first/version.txt。
Spark 把版本信息打印到标准错误。如果不用 2>&1 合并,文件就会是空的。除了版本号,还能看到它运行在哪个 Scala 和 Java 之上——Spark 4 要求 Java 17 或更高版本。
第一个作业——数 30 万条
以应用名 spk-first-count 创建 /root/spk/first/count.py,统计 /data/shop/orders.csv(带表头行)的行数,并以 {"rows": 정수}(占位符为整数)写入 /root/spk/first/out/count.json。用 spark-submit 运行。
SparkSession.builder.appName(...) 决定应用名。评分器会检查以该名字记录的事件日志是否完整关闭,以及文件中的数字是否与原始文件的行数相同。如果把表头行算作一行,就会多一个。
只堆叠转换算子,就不会启动作业
以应用名 spk-first-lazy 创建 /root/spk/first/lazy.py,用自己提供 schema 的 spark.read.csv 读取订单,堆叠 filter(或 where)和 withColumn,但不要调用任何行动算子。运行之后,该应用的事件日志里必须有 0 个作业。
不提供 schema 的话,Spark 为了确认表头行会启动一个小作业。如果像 schema="order_id STRING, ..." 这样传入 DDL 字符串,那个作业也会消失。打印 schema 没有问题——schema 驱动器早已知道,不会读取数据。
每个行动算子都会启动作业——用事件日志来数
以应用名 spk-first-actions 创建 /root/spk/first/actions.py,调用三次以上行动算子,并用其中一个把 100 行以 JSON 写入 /root/spk/first/out/sample。运行之后,在该应用的事件日志里数一数 SparkListenerJobStart 事件的个数,以一个整数写入 /root/spk/first/out/jobs.txt。
不要断定一个行动算子就是一个作业。也有写文件的行动算子或 schema 推断这样会产生两个以上作业的。所以要数的是日志,而不是代码。如果用同一个名字运行过多次,就要数最近的那份日志。
宽转换把阶段分开
以应用名 spk-first-agg 创建 /root/spk/first/agg.py,把按 status 统计的订单数,以带表头行的 CSV(列:status、count)写入 /root/spk/first/out/by_status。在该应用的事件日志里数一数 SparkListenerStageCompleted 事件的个数,以整数写入 /root/spk/first/out/stages.txt。
groupBy 需要把相同的键汇集到一个任务里,所以会产生 shuffle,而 shuffle 前后成为不同的阶段。评分器会检查你的 CSV 是否与从原始数据统计出的值相同、该应用中是否确实有用到 shuffle 的阶段、写下的阶段数是否与日志一致。
核心数决定分区数
创建 /root/spk/first/par.py,让它把第一个参数作为应用名,分别用 spark-submit --master 'local[1]' par.py spk-first-par1 和 spark-submit --master 'local[2]' par.py spk-first-par2 运行两次,把 {"master": 문자열, "default_parallelism": 정수, "partitions": 정수}(占位符依次为字符串、整数、整数)写入 /root/spk/first/out/spk-first-par1.json 和 /root/spk/first/out/spk-first-par2.json。partitions 是提供 schema 读取的订单 DataFrame 的 rdd.getNumPartitions()。
文件被切成几块,并不只由 spark.sql.files.maxPartitionBytes(默认 128MB)决定。如果有多个核心,当总大小除以核心数的值更小时,就按这个大小来切——这是为了不让核心闲着的规则。看看一个 16MB 的文件,在一个核心和两个核心下有什么不同。
失败也是应用——读取错误条件名
以应用名 spk-first-fail 创建 /root/spk/first/fail.py,读取不存在的路径 /data/shop/order.csv,捕获异常,把 getCondition() 返回的错误条件名写到 /root/spk/first/out/error.txt 的第一行。应用必须用 spark.stop() 正常结束。
不存在的路径不是在调用行动算子时,而是在 read 的那一刻暴露的。因为文件列表是在制定计划时确认的。Spark 4 的错误,除了给人读的句子之外,还有给机器读的名字(错误条件)。告警和重试规则最好基于这个名字,而不是句子来设置。
用数字留下看到了什么
在 /root/spk/first/report.md 中写出 ## 잡과 스테이지(韩文,意为“作业与阶段”)、## 게으른 실행(韩文,意为“惰性执行”)、## 파티션(韩文,意为“分区”)三节。第一节放入第 4 步数出的作业数,第三节放入第 6 步的两个分区数,都用数字写出。
不要只抄数字,要各配一两句话说明为什么是这个数字。三个行动算子对应几个作业,只堆叠转换算子的应用为什么是 0 个,核心数如何改变了分区数,这就是本实验看到的全部。