PySpark的collect_list和collect_set底层用groupByKey还是reduceByKey?
PySpark
collect_list与collect_set底层聚合逻辑说明 核心结论
collect_list和collect_set底层均采用与groupByKey一致的全量shuffle逻辑,不会使用reduceByKey的map端预聚合优化,执行时需要将同一个key的所有对应数据通过shuffle传输到下游reduce节点后再完成聚合。
为什么不使用reduceByKey模式
reduceByKey的核心优化前提是聚合操作同时满足结合律和交换律,可以在map端先做局部聚合大幅减少需要shuffle的数据量,但这两个聚合函数不满足优化前提:
collect_list需要保留同一个key下元素的出现顺序,必须全量收集所有元素后才能得到最终结果,map端预聚合没有任何收益,反而会增加序列化、反序列化的额外开销collect_set虽然需要对元素去重,但也需要收集全量元素后完成全局去重,即使部分高版本Spark支持map端局部去重的小优化,本质上还是需要将所有局部去重后的元素传输到reduce节点,和reduceByKey的聚合逻辑有本质区别
验证方式
你可以通过在聚合代码后调用explain()方法打印物理执行计划,查看shuffle阶段的输出数据量,会发现map端的输出和原始对应key的数据量基本一致,没有出现sum、count这类支持预聚合的算子常见的map端数据量缩减情况。
内容的提问来源于stack exchange,提问作者aishik roy chaudhury
相关产品推荐
相关产品推荐

