Apache Spark中requiredChildDistribution与outputPartitioning有什么区别?
Spark物理执行计划算子核心属性差异说明
Spark的物理算子四个属性本质可以分成两类:算子自身输出特征、算子对子节点的输入要求,两两对应,不存在等价关系。
outputPartitioning 与 requiredChildDistribution 的差异
outputPartitioning是当前算子的输出特征声明:用来标记当前算子计算完成后,输出的数据集采用的分区规则,常见的包括哈希分区、范围分区、单分区、随机分区等。下游算子可以直接读取该属性,判断是否需要额外执行shuffle来适配自身的执行要求。requiredChildDistribution是当前算子的输入要求声明:用来标记当前算子要正常执行,其所有子节点的输出数据分布必须满足的规则。比如普通哈希聚合算子就要求子节点输出必须按聚合键做哈希分布,如果子节点的outputPartitioning不满足该要求,Spark会自动在子节点和当前算子之间插入shuffle阶段调整数据分布。
实际场景示例:执行
df.groupBy("user_id").agg(sum("pay_amount"))时,聚合算子的requiredChildDistribution要求输入按user_id哈希分布。如果上游算子的输出已经是按user_id哈希分区的状态,Spark就会直接跳过shuffle步骤,否则必须触发shuffle来重分区。
outputOrdering 与 requiredChildOrdering 的差异
outputOrdering是当前算子的输出特征声明:用来标记当前算子输出的数据集,在单个分区内部的排序规则,包括排序字段、升降序、空值排列规则等。requiredChildOrdering是当前算子的输入要求声明:用来标记当前算子要正常执行,其所有子节点输出的单分区内数据必须满足的排序规则。比如SortMergeJoin算子就要求两个子节点的输出都必须在各自分区内按Join键排序,如果子节点的outputOrdering不满足该要求,Spark会自动在子节点和当前算子之间插入排序算子。
这两类属性是Spark Catalyst优化器做物理执行优化的核心判断依据,我们常说的「消除冗余shuffle」「消除冗余排序」优化,本质就是优化器判断上游算子的输出特征,已经完全匹配下游算子的输入要求,因此跳过了多余的计算步骤。
内容的提问来源于stack exchange,提问作者Arjunlal M.A
相关产品推荐
相关产品推荐

