Apache Spark — 慢作业的答案在执行计划和事件日志里
DataFrame 不是结果而是配方,所以每次都重新计算
一句话总结
DataFrame 不是计算好的结果,而是从原始数据走到结果的谱系(lineage),所以每次调用行动算子,都要从原始数据重新计算。缓存是把第一次计算的结果攥在手里,检查点则是把结果写成文件并切断谱系。
为什么同一个计算要做两次
假设读取 CSV、清洗、连接得到一个 DataFrame,先调用 count(),接着再把同一个 DataFrame 写成文件。光看代码,清洗只有一次。实际上发生了两次。看事件日志,两个作业都是从 CSV 扫描开始的。
原因是转换是惰性的。filter 或 join 只是往计划里加一行,什么也不计算。调用行动算子时,Spark 把那个计划从原始数据一直执行到最后,把结果交给驱动器或写成文件之后,就把中间结果丢掉。下一个行动算子又从原始数据开始。这不是缺陷,而是设计。把中间结果全部留着,内存很快就满了,而有了谱系,丢掉的碎片随时可以重新生成。
所以只挑出要多次使用的中间结果来留存,是人的职责。这就是缓存。
工作原理
cache() 也是惰性的。调用的那一刻,只留下“第一次计算这个结果时把它留下来”的标记。第一个行动算子运行时,各个分区一边计算一边被保存,之后的行动算子在计划里不再从原始数据扫描开始,而是从 InMemoryTableScan 开始。正如 RDD 编程指南所说,它在第一次计算时被保存,之后的行动算子会重复使用它。
from pyspark import StorageLevel
clean = raw.filter("qty > 0").join(products, "product_id")
clean.cache() # 표시만 한다 — 아직 아무것도 저장되지 않았다
clean.count() # 첫 행동: 계산하면서 저장
clean.write.parquet("/tmp/out") # 계획이 InMemoryTableScan 에서 시작한다
clean.unpersist() # 다음 행동은 다시 원본부터
slim = raw.select("order_id", "qty").persist(StorageLevel.DISK_ONLY)
存储级别决定把它放在哪里。容易混淆的是默认值。RDD 的 cache() 是只放在内存里的 MEMORY_ONLY。而 DataFrame 的 cache() 是 MEMORY_AND_DISK_DESER,装不进内存的分区放在磁盘上。文档写道,这个默认值是在 3.0 中为了与 Scala 一致而改的。想要其他级别,就传给 persist(),但已经确定了存储级别的 DataFrame,不能再给它新的级别。必须先 unpersist()。
DataFrame 缓存不是按原来的行,而是按列式格式放在内存里的。性能调优文档说明,缓存的表只读取需要的列,并自动调整压缩。
关于选哪个级别,RDD 指南的建议很简短。如果能舒服地装进内存,就用默认值,除非计算很贵,或者过滤掉的量很大,否则不要溢写到磁盘。因为重新计算可能和从磁盘读取一样快。如果原始数据是一个本地 CSV,转换又很轻,DISK_ONLY 缓存几乎没有好处。
缓存所在的地方——统一内存
缓存不是放在免费的空间里。根据调优指南,shuffle、连接、排序、聚合使用的执行内存,与缓存使用的存储内存,共用同一个区域。这个区域的大小是 spark.memory.fraction(默认 0.6),它是相对于 JVM 堆减去 300MiB 之后剩余部分的比例。在这个区域里,spark.memory.storageFraction(默认 0.5)那一份,是执行内存无法挤占缓存的份额。
按本实验中驱动器内存 1g 来估算,(1024 - 300) × 0.6,也就是 400MiB 出头,由执行和缓存共用。缓存了大表,连接和排序可用的位置就相应减少;缓存溢出的话,正如 RDD 指南所说,从最久没用的分区开始(LRU)被挤出去。被挤出去的分区,下次需要时会沿着谱系重新计算。所以缓存不会给出错误的答案,只会变慢。
unpersist 与检查点
用完的缓存,要用 unpersist() 放掉。RDD 指南写道,这个调用默认不会等待。想等到资源真正释放,就传入 blocking 参数。放掉之后的行动算子,在计划里又会从原始数据扫描开始。还有一点,DataFrame 缓存由集群的所有会话共享。如果别的会话已经把同一个计划缓存了,我的查询也会用到它。
缓存保留谱系。丢失的分区之所以能重新生成,靠的就是这个。但是像迭代算法那样,在同一个 DataFrame 上叠加几十次转换,谱系本身就成了问题。检查点文档说明,这种情况下计划可能呈指数级增长,而检查点会截断逻辑计划。它把结果写成文件放到由 SparkContext.setCheckpointDir() 或 spark.checkpoint.dir 指定的目录里,之后的计划不再从原始数据开始,而是从那个文件开始。计划里会出现读取已有 RDD 的节点,来代替原始数据扫描。默认是即时(eager)执行,所以调用的那一刻就会运行作业。
localCheckpoint() 不写文件,而是保存到执行器的缓存体系里。它很快,但正如文档自己所说,并不可靠——丢失执行器时,谱系也已经被截断,没有办法重新生成。
SQL 里也有同样的工具。CACHE TABLE 与 DataFrame 的 cache() 不同,默认是即时的,必须加上 LAZY,才会在第一次使用时缓存。不单独指定存储级别,就是 MEMORY_AND_DISK。
在现场相遇的样子
第一,把只用一次的东西缓存了。对只有一个行动算子的 DataFrame 做缓存,只会多花保存的成本。缓存只放在被使用两次以上的分叉处。
第二,把缓存攥着然后忘了。在长笔记本或服务里,缓存一旦堆积,执行内存就会减少,shuffle 开始溢写到磁盘。用完就放掉。
第三,以为缓存了,计划里却没有。如果对在 cache() 之后又加了新转换的另一个 DataFrame 调用行动算子,那么在那个 DataFrame 的计划里,只有被缓存的那一部分会被重用。缓存是否真的被用到,要在行动算子之后,到 explain() 里找 InMemoryTableScan 来确认。Web UI 的 Storage 选项卡也要等到物化之后才会显示缓存。
实际工作中真正重要的事
- DataFrame 是食谱。每个行动算子都要从原始数据重新计算。
- cache() 是惰性的。在第一个行动算子运行时才保存。
- DataFrame 的默认存储级别是 MEMORY_AND_DISK_DESER,RDD 是 MEMORY_ONLY。
- 缓存与执行内存共用同一个区域。用完就 unpersist。
- 谱系太长就用检查点切断。缓存不会切断谱系。
下一项实验要做什么
不用缓存调用两个行动算子,通过事件日志和计划确认两个作业都要重新读取 CSV。在 cache() 之后,查看第二个行动算子的计划是否从 InMemoryTableScan 开始,以及存储级别,并用 persist 指定只放在磁盘上的级别。确认在 unpersist 前后缓存状态会改变,把叠加了十次转换的 DataFrame 用检查点切断,看是否得到同样的答案,再用 SQL 的 CACHE TABLE 做同样的事。最后,在开启了块更新记录的应用的事件日志里,把缓存块的大小相加,用数字测出缓存占用了多少内存。