Spark引入GROUP BY后查询变慢及执行计划解读求助
问题解答
一、关于查询性能差异的理解是否正确?
完全正确,再补充细节帮你加深认知:
- 第一个查询:
display(joined_df.limit(100))中的limit会触发Limit下推优化,Spark的自适应执行(AdaptiveSparkPlan)在获取到100条符合条件的结果后,会直接终止后续的表扫描和关联操作,根本不需要全量处理1400万行的items表,所以能快速完成。从执行计划也能看到,CollectLimit在最上层,所有关联逻辑都是为了获取这100行数据服务。 - 第二个查询:GROUP BY搭配MIN聚合必须处理全量关联后的数据集——首先要完成三张表的完整关联,接着按order_id排序做局部聚合,再通过Exchange shuffle数据到全局聚合阶段,最后才能得到结果。再加上items表有1400万行,单节点(4核32GB)的计算和内存负载远超能力范围,导致查询无法完成。
二、执行计划中的+符号含义
这些+是层级缩进标记,用来清晰展示算子之间的依赖关系:
- 每一层
+代表当前算子是上一层算子的子节点,即上一层算子要依赖当前算子的输出结果 - 比如
CollectLimit (13)下面的+- BroadcastHashJoin Inner BuildLeft (12),意思是BroadcastHashJoin的输出会传给CollectLimit做处理 - 它和stage、job没有直接对应关系
三、如何判断新Job或Stage的创建
1. Job的判断规则
Spark里每个Action算子(比如display、collect、write等)会触发一个独立的Job。你的两个查询各自对应一个Action,所以分别是1个Job。
2. Stage的判断规则
Stage的划分边界是Shuffle操作(Exchange算子),只要出现Exchange,就会划分出新的Stage:
- 执行计划里的
ShuffleQueryStage、ResultQueryStage就是具体的Stage实例 - 比如第一个执行计划中的
ShuffleQueryStage (9)、ShuffleQueryStage (4),都是由前面的Exchange算子触发的Stage划分 - 简单总结:从数据源扫描开始,到第一个Exchange之前是一个Stage;两个Exchange之间是一个Stage;最后一个Exchange到输出环节是一个Stage
内容的提问来源于stack exchange,提问作者ng.newbie
相关产品推荐
相关产品推荐

