Apache Spark — 慢作业的答案在执行计划和事件日志里
避免把同一计算做两次——缓存在执行计划中的样子
目标
对同一个 DataFrame 调用两次行动算子时,通过计划确认没有缓存就会从原始数据重新计算,并查看 cache、persist、unpersist、检查点、CACHE TABLE 在计划和事件日志里是怎么显示的。还要用块更新事件测出缓存实际占用了多少内存。
为什么重要
DataFrame 不是结果,而是制作的方法(谱系)。每次调用行动算子,Spark 都会把这个方法从头再执行一遍。把同一个中间结果用两次的代码,会把原始数据读两次、把连接做两次——既没有错误,也没有警告。
cache() 是一个标记,表示把这个中间结果留在执行器里。因为只是标记,所以由第一个行动算子来填充,之后的行动算子在计划里不再读取原始数据,而是用 InMemoryTableScan 读取缓存。存储级别(persist)是选择内存、磁盘、是否序列化,决定内存不足时放弃什么。用完的缓存,必须用 unpersist 放下,才能成为其他作业的内存。
谱系变得太长(反复计算)的话,计划本身就会变大,驱动器会变慢。检查点把结果写成文件并切断谱系,让之后的计划变成“读取这个文件”这一行。
步骤
- 在 /root/spk/cache/common.py 中放入给已完成支付订单加上金额的函数,以及调用两次行动算子(
count、sum(amount))的函数,再用 /root/spk/cache/none.py(应用spk-cache-none),在不缓存的情况下把两个行动算子的结果写入 /root/spk/cache/out/none.json。 - 在 /root/spk/cache/cached.py(应用
spk-cache-on)中,给同一个 DataFrame 加上cache(),把两个行动算子的结果写入 /root/spk/cache/out/cached.json,把存储级别(str(df.storageLevel))写入 /root/spk/cache/out/cached_level.txt。 - 在 /root/spk/cache/disk.py(应用
spk-cache-disk)中,用persist(StorageLevel.DISK_ONLY)做同样的事,把结果写入 /root/spk/cache/out/disk.json,把存储级别写入 /root/spk/cache/out/level.txt。 - 在 /root/spk/cache/unpersist.py(应用
spk-cache-un)中,按缓存 → 行动算子 →unpersist()→ 行动算子的顺序运行,把是否被缓存(df.is_cached)在前后的值,以{"before": 참거짓, "after": 참거짓}(占位符均为布尔值)写入 /root/spk/cache/out/unpersist.json。 - 在 /root/spk/cache/ckpt.py(应用
spk-cache-ckpt)中,把检查点文件夹设为 /root/spk/cache/ckpt,在给金额依次加上 1–10 的十次withColumn之后执行checkpoint(),用结果把按渠道的金额之和以 Parquet(channel、amount)写入 /root/spk/cache/out/ckpt_result。 - 在 /root/spk/cache/cache_sql.py(应用
spk-cache-sql)中创建临时视图paid,执行CACHE TABLE paid_cached AS SELECT channel, amount FROM paid之后求出按渠道的合计,以{"is_cached": 참거짓, "by_channel": {채널: 합}}(占位符依次为布尔值、渠道、合计)写入 /root/spk/cache/out/cache_table.json。 - 用 /root/spk/cache/size.py(应用
spk-cache-size,spark.eventLog.logBlockUpdates.enabled=true)把缓存填满,然后把该日志里rdd_块的大小相加,以{"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수}(占位符均为整数)写入 /root/spk/cache/out/cache_size.json。 - 在 /root/spk/cache/report.md 中以
## 다시 계산하는 비용(韩文,意为“重新计算的成本”)、## 저장 수준(韩文,意为“存储级别”)、## 캐시의 크기(韩文,意为“缓存的大小”)三节写成报告。第二节放入第 2、3 步的存储级别字符串,第三节放入第 7 步的块数和内存字节数。
参考
- 脚本放在
/root/spk/cache里并在那里运行(from common import …)。原始数据是/data/shop/orders.csv、/data/shop/products.csv,金额 = qty × price。 - 读取缓存的那次运行的计划里,会看到
InMemoryTableScan以及它下面的InMemoryRelation … StorageLevel(…)。原始数据扫描(Scan csv)仍留在 InMemoryRelation 的说明里,但不会被执行。 - 块更新事件(
SparkListenerBlockUpdated)中,同一个块可能出现多次。每个块 ID 只保留最后一条再相加。 - 常见错误:只做了
cache()而没调用行动算子,导致缓存是空的;缓存用完不放掉;在一个应用里把原始和派生的 DataFrame 都缓存了,分不清读的是哪个缓存。 - 官方文档:RDD Persistence · Caching Data · CACHE TABLE · Configuration — spark.eventLog.logBlockUpdates.enabled
没有缓存,两个行动算子——原始数据被读了两次
在 /root/spk/cache/common.py 中放入:把已完成支付订单与商品连接并加上 amount(qty×price)的函数,以及对 DataFrame 调用 count() 和 sum(amount) 两个行动算子、并把 {"count": 정수, "sum": 정수}(占位符均为整数)写入文件的函数。以应用名 spk-cache-none 创建 /root/spk/cache/none.py,在不缓存的情况下把结果写入 /root/spk/cache/out/none.json。
有两个行动算子,SQL 执行也就有两次,两次的计划最下面都有 Scan csv。这意味着同一个文件被读了两次,同样的连接也做了两次。评分器会检查没有缓存、读取了原始数据的运行是否有两次以上。
cache()——第一个行动算子来填充,下一个来读取
以应用名 spk-cache-on 创建 /root/spk/cache/cached.py,对第 1 步的 DataFrame 加上 cache(),把同样的两个行动算子的结果写入 /root/spk/cache/out/cached.json,把 str(df.storageLevel) 写入 /root/spk/cache/out/cached_level.txt。
第二次运行的计划里会出现 InMemoryTableScan。还要留意,DataFrame 的默认存储级别与 RDD 的 cache 不同——内存不足时会溢写到磁盘。结果数字必须与第 1 步相同。
persist(DISK_ONLY)——选择存储级别
以应用名 spk-cache-disk 创建 /root/spk/cache/disk.py,把第 1 步的 DataFrame 用 persist(StorageLevel.DISK_ONLY) 保存,把同样的两个行动算子的结果写入 /root/spk/cache/out/disk.json,把 str(df.storageLevel) 写入 /root/spk/cache/out/level.txt。
只放在磁盘上,可以节省执行器内存,但每次读取都需要反序列化。看看计划的 InMemoryRelation 那一行里打印的 StorageLevel 发生了什么变化。名字里带 InMemory,也可能是在磁盘上。
unpersist——放下缓存,就要重新计算
以应用名 spk-cache-un 创建 /root/spk/cache/unpersist.py,把第 1 步的 DataFrame 按 cache() → count() → unpersist() → sum(amount) 的顺序运行,把 unpersist() 前后的 df.is_cached,以 {"before": 참거짓, "after": 참거짓}(占位符均为布尔值)写入 /root/spk/cache/out/unpersist.json。
缓存不是免费的。它占用执行器内存,减少其他作业可用的位置。用完的缓存要放下。评分器会通过日志检查,读取了缓存的运行之后,是否接着有不读取缓存的运行。
用检查点切断谱系
以应用名 spk-cache-ckpt 创建 /root/spk/cache/ckpt.py,把检查点文件夹设为 /root/spk/cache/ckpt,对第 1 步的 DataFrame,给 amount 依次加上 1、2、…、10 的十次 withColumn 之后,执行 checkpoint(),用结果把按渠道的 amount 之和以 Parquet(channel、amount)写入 /root/spk/cache/out/ckpt_result。
checkpoint() 会当场把结果写入检查点文件夹,它返回的 DataFrame 的计划不再是原始数据,而是读取那个文件的一行(Scan ExistingRDD)。缓存保留谱系(丢失了就重新计算),检查点切断谱系(文件就是原始数据)。
SQL 的 CACHE TABLE 不是惰性的
以应用名 spk-cache-sql 创建 /root/spk/cache/cache_sql.py,把第 1 步的 DataFrame 注册成临时视图 paid,执行 CACHE TABLE paid_cached AS SELECT channel, amount FROM paid,然后从 paid_cached 求出按渠道的 amount 之和,以 {"is_cached": spark.catalog.isCached("paid_cached"), "by_channel": {채널: 합}}(占位符依次为渠道、合计)写入 /root/spk/cache/out/cache_table.json。
DataFrame 的 cache() 只是打标记,而 CACHE TABLE 当场就填充(如果是 CACHE LAZY TABLE 则推迟)。所以 CACHE TABLE 语句本身就会启动作业。读取有名字的缓存的计划里,会打印 Scan In-memory table paid_cached。评分器会检查有没有这样的运行,以及合计。
测出缓存实际占用的内存
以应用名 spk-cache-size、设置 spark.eventLog.logBlockUpdates.enabled=true 创建 /root/spk/cache/size.py,对第 1 步的 DataFrame 执行 cache(),并用 count() 把它填满。在该应用的日志里,从 SparkListenerBlockUpdated 中挑出块 ID 以 rdd_ 开头的,每个块只保留最后一条,把个数以及 Memory Size、Disk Size 之和,以 {"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수}(占位符均为整数)写入 /root/spk/cache/out/cache_size.json。
缓存块是每个分区一个(rdd_<RDD 번호>_<파티션>,占位符依次为 RDD 编号、分区)。在内存里展开(反序列化)的缓存,可能比原始 CSV 大,也可能小——取决于选择的列和类型。在增加缓存之前先看这个数字的习惯,会保住执行器内存。
留下缓存该放哪里、该放掉哪里
在 /root/spk/cache/report.md 中写出 ## 다시 계산하는 비용(韩文,意为“重新计算的成本”)、## 저장 수준(韩文,意为“存储级别”)、## 캐시의 크기(韩文,意为“缓存的大小”)三节。第二节放入第 2、3 步的存储级别字符串,第三节放入第 7 步的块数和内存字节数。
第一节写没有缓存时把什么做了两次,第二节写两种存储级别在交换什么,第三节写这些内存占执行器堆(1g)的百分之几。