Apache Spark — 慢作业的答案在执行计划和事件日志里
洗牌是最昂贵的一步,200 个分区是一个不了解数据就定下的数字
一句话总结
shuffle 是为了把相同的键汇集到一处,而让所有分区向所有分区发送数据,它要同时付出磁盘写入、序列化和网络的代价。shuffle 之后的分区数默认是 200,这是与数据大小无关的固定数字。AQE 会在运行过程中合并较小的分区,但究竟合并了什么,以及 repartition 与 coalesce 的区别,必须亲自确认。
为什么 shuffle 是个问题
filter 或 withColumn 接收一个分区,输出一个分区。不需要看相邻的分区,所以在一个任务里就能完成。这样的叫窄转换。groupBy、join、distinct、orderBy 则不同。不知道相同的键散落在哪些分区里,所以要计算,就必须翻遍所有分区,把相同的键汇集到一起。
RDD 指南中的 shuffle 一节把它称为在分区之间重新分配数据的 Spark 机制,并列出了它开销很大的三个原因:磁盘输入输出、数据序列化、网络输入输出。准备 shuffle 的一方叫 Map 任务,汇集后计算的一方叫 Reduce 任务,文档补充说,这两个名字来自 MapReduce,与 Spark 的 map、reduce 运算并没有直接关系。
工作方式是这样的。Map 一侧的任务先把结果收集在内存里,溢出时按目标分区的顺序排序后写成文件。Reduce 一侧的任务从所有 Map 任务写出的文件中,只挑出属于自己的那一份块来读取。文档写道,shuffle 会在磁盘上产生大量中间文件,并且为了在重新计算谱系时能用上,会把这些文件一直保留到引用消失。所以 shuffle 很多的长作业会相当耗磁盘。
工作原理——200 这个数字
shuffle 之后要分成几个分区,必须有人来定。在 DataFrame 和 SQL 里,这个值是性能调优文档中的 spark.sql.shuffle.partitions,它是连接或聚合做 shuffle 时使用的分区数,默认值是 200。
200 不是看着数据选出来的数字。对 1TB 来说太少,一个分区会有 5GB 左右,内存溢出,数据会溢写到磁盘。对 30 万行来说又太多,每个分区里只有几百行,启动和回收 200 个任务的固定开销超过了工作本身。如果把结果原样写成文件,每个写入任务都会写一个文件,于是会涌出大量小文件。
spark.conf.set("spark.sql.adaptive.enabled", "false") # AQE 를 끄고 날것을 본다
by_day = orders.groupBy("day").agg(F.sum("qty"))
by_day.write.mode("overwrite").parquet("/root/spk/out/by_day") # 셔플 뒤 태스크 200개
spark.conf.set("spark.sql.shuffle.partitions", "8") # 자료에 맞춰 줄인다
AQE 的分区合并
手工调整数字之所以困难,是因为 shuffle 之后的大小,在 shuffle 之前是无从得知的。合并 shuffle 之后的分区一节解决了这个问题。当 AQE 和 spark.sql.adaptive.coalescePartitions.enabled 都开启时(默认都是 true),会根据 Map 一侧的输出统计信息,合并相邻的较小 shuffle 分区。如果不指定起始分区数 initialPartitionNum,它就等于 spark.sql.shuffle.partitions。
合并到什么程度,还有细节。目标大小 advisoryPartitionSizeInBytes 的默认值是 64MB,但由于 parallelismFirst 默认为 true,所以会无视这个目标大小,只遵守最小大小 minPartitionSize(默认 1MB),把并行度设到最大。文档建议在繁忙的集群里把它设为 false,以免出现过多的小任务。也就是说,默认设置下的 AQE 并不是“聚成 64MB”。合并后的结果会在最终计划中显示为 AQEShuffleRead ... coalesced。
repartition 与 coalesce
直接改变分区数的方法有两种,成本完全不同。repartition 会创建一个经过哈希分区的新 DataFrame。它要把所有行重新分配,所以会产生 shuffle。作为代价,分区大小会变得均匀,而且如果指定了列,相同的值会汇集到同一个分区。
coalesce 是窄依赖。用文档的例子,从 1000 个减少到 100 个时,不需要 shuffle,一个新分区负责 10 个原有分区。如果要求增加,则停留在当前的数量。不过文档警告说,像 coalesce(1) 这样的急剧缩减,可能让计算本身只在很少的节点上运行。因为没有 shuffle,所以前面的阶段都会被合并成一个任务。这时,即使要付出一次 shuffle 的代价,repartition 也能让前面的阶段并行运行。
RDD 指南把 coalesce 也列进了可能引发 shuffle 的运算里,这是因为 RDD 的 coalesce 把是否 shuffle 当作参数来接收。DataFrame 的 coalesce 如上所述,没有 shuffle。在 SQL 里,用分区提示 COALESCE、REPARTITION、REBALANCE 来做同样的事,文档也把它介绍为减少输出文件数的工具。
在现场相遇的样子
小数据却有 200 个文件。用开发数据运行的作业,在结果文件夹里生成了一大堆只有几十字节的文件。只看一个设置就能找到原因。开启了 AQE 的话会减少很多,但在写入之前用 coalesce 来确定文件数更可靠。
shuffle 的量在日志和 UI 里。Web UI 的 Stages 选项卡会按阶段显示 Shuffle read 和 Shuffle write 的字节数和记录数,事件日志的任务结束事件里也有同样的指标。不要凭“shuffle 很大”的感觉,而要用数字说话。
只有窄转换的作业只有一个阶段。只做读取、过滤、写入的作业没有 shuffle,所以只有一个阶段。如果阶段不止一个,就是某处有 shuffle 的证据。
实际工作中真正重要的事
- shuffle 要同时付出磁盘、序列化、网络的代价。阶段边界就是 shuffle。
- shuffle.partitions 的 200 是不了解数据的默认值。要按数据大小来定,或者交给 AQE,但要确认结果。
- 默认的 AQE 不会聚成 64MB。parallelismFirst 为 true,所以优先保证并行度。
- 减少用 coalesce,均匀划分用 repartition。前者没有 shuffle,后者有 shuffle。
- coalesce(1) 会把前面的阶段也变成一个任务。大数据要用 repartition。
- shuffle 的字节数要在 UI 和事件日志中用数字读取。
下一项实验要做什么
关闭 AQE,汇总按客户的销售额并写出,数一数 200 个 shuffle 分区产生的结果文件数,再把 shuffle 分区减到 8,比较差异。把设置恢复成默认值,从事件日志里数一数 AQE 合并之后实际读取了 shuffle 的任务数;对点击数据分别施加 repartition 和 coalesce,比较结果文件数。直接从日志里汇总 shuffle 写出的字节数,并确认只有窄转换的作业没有 shuffle,再整理成报告。