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

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

DataFrame 不是结果而是配方,所以每次都重新计算

在 TT Lab 中继续学习

一句话总结

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 选项卡也要等到物化之后才会显示缓存。

实际工作中真正重要的事

下一项实验要做什么

不用缓存调用两个行动算子,通过事件日志和计划确认两个作业都要重新读取 CSV。在 cache() 之后,查看第二个行动算子的计划是否从 InMemoryTableScan 开始,以及存储级别,并用 persist 指定只放在磁盘上的级别。确认在 unpersist 前后缓存状态会改变,把叠加了十次转换的 DataFrame 用检查点切断,看是否得到同样的答案,再用 SQL 的 CACHE TABLE 做同样的事。最后,在开启了块更新记录的应用的事件日志里,把缓存块的大小相加,用数字测出缓存占用了多少内存。