Spark是否优化高阶链式转换?复用sorted_array会重复排序吗?
Spark是否会优化高阶链式转换?
在以下示例中,若创建N个使用sorted_array的列,Spark会对原始数组执行N次排序,还是会通过WholeStageCodegen进行优化?
from pyspark.sql import functions as F _data = [ ([1, 10, 12, 2, 4], ), ([1, 100, 4, 1, 5], ) ] df = spark.createDataFrame(_data, ['arr_col']) sorted_array = F.sort_array('arr_col') cols = [ F.slice(sorted_array, 1, 1).alias('c1'), F.slice(sorted_array, 1, 2).alias('c2'), F.slice(sorted_array, 1, 3).alias('c3') ] df2 = df.select(cols) df2.explain()
执行后输出的物理计划:
# == Physical Plan == # *(1) Project [ # slice(sort_array(arr_col#0, true), 1, 1) AS c1#43, # slice(sort_array(arr_col#0, true), 1, 2) AS c2#44, # slice(sort_array(arr_col#0, true), 1, 3) AS c3#45 # ] # +- *(1) Scan ExistingRDD[arr_col#0]
从物理计划可以直接看到,当前Spark不会自动优化重复的sort_array计算,它会对原始数组执行N次排序(这里N=3)。
Spark的逻辑优化器默认不会复用相同表达式的计算结果,哪怕表达式完全一致。而WholeStageCodegen的作用是把整个Stage的算子逻辑合并成一段原生代码执行,减少JVM上下文切换的开销,但它不会消除重复的计算逻辑,只是提升单次计算的执行效率。
如果想要避免重复排序,需要显式将排序后的数组作为中间列保存,再基于这个列生成后续的slice列,示例代码如下:
from pyspark.sql import functions as F _data = [ ([1, 10, 12, 2, 4], ), ([1, 100, 4, 1, 5], ) ] df = spark.createDataFrame(_data, ['arr_col']) # 先添加排序后的中间列 df_with_sorted = df.withColumn('sorted_arr', F.sort_array('arr_col')) cols = [ F.slice('sorted_arr', 1, 1).alias('c1'), F.slice('sorted_arr', 1, 2).alias('c2'), F.slice('sorted_arr', 1, 3).alias('c3') ] df2 = df_with_sorted.select(cols) df2.explain()
对应的物理计划会变成:
# == Physical Plan == # *(1) Project [ # slice(sorted_arr#46, 1, 1) AS c1#49, # slice(sorted_arr#46, 1, 2) AS c2#50, # slice(sorted_arr#46, 1, 3) AS c3#51 # ] # +- *(1) Project [sort_array(arr_col#0, true) AS sorted_arr#46] # +- *(1) Scan ExistingRDD[arr_col#0]
此时sort_array只会执行一次,后续的slice操作都基于这个已排序的数组进行。
内容的提问来源于stack exchange,提问作者boyangeor
相关产品推荐
相关产品推荐

