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

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

避免把同一计算做两次——缓存在执行计划中的样子

在 TT Lab 中继续学习

目标

对同一个 DataFrame 调用两次行动算子时,通过计划确认没有缓存就会从原始数据重新计算,并查看 cache、persist、unpersist、检查点、CACHE TABLE 在计划和事件日志里是怎么显示的。还要用块更新事件测出缓存实际占用了多少内存。

为什么重要

DataFrame 不是结果,而是制作的方法(谱系)。每次调用行动算子,Spark 都会把这个方法从头再执行一遍。把同一个中间结果用两次的代码,会把原始数据读两次、把连接做两次——既没有错误,也没有警告。 cache() 是一个标记,表示把这个中间结果留在执行器里。因为只是标记,所以由第一个行动算子来填充,之后的行动算子在计划里不再读取原始数据,而是用 InMemoryTableScan 读取缓存。存储级别(persist)是选择内存、磁盘、是否序列化,决定内存不足时放弃什么。用完的缓存,必须用 unpersist 放下,才能成为其他作业的内存。 谱系变得太长(反复计算)的话,计划本身就会变大,驱动器会变慢。检查点把结果写成文件并切断谱系,让之后的计划变成“读取这个文件”这一行。

步骤

  1. 在 /root/spk/cache/common.py 中放入给已完成支付订单加上金额的函数,以及调用两次行动算子(count、sum(amount))的函数,再用 /root/spk/cache/none.py(应用 spk-cache-none),在不缓存的情况下把两个行动算子的结果写入 /root/spk/cache/out/none.json。
  2. 在 /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。
  3. 在 /root/spk/cache/disk.py(应用 spk-cache-disk)中,用 persist(StorageLevel.DISK_ONLY) 做同样的事,把结果写入 /root/spk/cache/out/disk.json,把存储级别写入 /root/spk/cache/out/level.txt。
  4. 在 /root/spk/cache/unpersist.py(应用 spk-cache-un)中,按缓存 → 行动算子 → unpersist() → 行动算子的顺序运行,把是否被缓存(df.is_cached)在前后的值,以 {"before": 참거짓, "after": 참거짓}(占位符均为布尔值)写入 /root/spk/cache/out/unpersist.json。
  5. 在 /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。
  6. 在 /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。
  7. 用 /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。
  8. 在 /root/spk/cache/report.md 中以 ## 다시 계산하는 비용(韩文,意为“重新计算的成本”)、## 저장 수준(韩文,意为“存储级别”)、## 캐시의 크기(韩文,意为“缓存的大小”)三节写成报告。第二节放入第 2、3 步的存储级别字符串,第三节放入第 7 步的块数和内存字节数。

参考

没有缓存,两个行动算子——原始数据被读了两次

在 /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)的百分之几。