You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 00:04:59