Azure Databricks中SQL查询因JOIN字段选择导致执行时长差异的问题
Azure Databricks SQL查询性能差异的底层原理分析
先看你遇到的两个查询:
慢查询(耗时超1小时)
SELECT DISTINCT t.TradeAmount FROM crm.NewTransaction t LEFT OUTER JOIN drd.PriceHistoric PPH ON t.TradeDate = PPH.PriceDate
快查询(不到1秒)
SELECT DISTINCT t.TradeAmount ,t.TradeDate FROM crm.NewTransaction t LEFT OUTER JOIN drd.PMSPriceHistoric PPH ON t.TradeDate = PPH.PriceDate
两者性能天差地别的核心原因在于Databricks Catalyst优化器的**连接消除(Join Elimination)**规则触发与否,以及后续计算开销的差异:
连接消除的触发逻辑
第二个查询的输出字段包含了Join关联键t.TradeDate,优化器可以明确判断:Left Outer Join不会改变t.TradeAmount和t.TradeDate的取值——即使PPH表中没有匹配TradeDate的记录,这两个字段依然是原表的数值,且最终的DISTINCT操作只针对这两个字段,完全可以直接从crm.NewTransaction表读取数据并去重,不需要执行实际的Join操作。这一步优化直接跳过了两张表关联的巨大开销。无法消除连接的性能瓶颈
第一个查询只输出t.TradeAmount,优化器无法跳过Join:因为PPH表中同一个TradeDate可能有多条记录,Join后会让同一个t.TradeAmount对应多条重复记录,最终的DISTINCT必须基于Join后的全量数据计算。这个过程需要完成:- 两张表按照
TradeDate做全量关联,生成远大于原表的中间数据集; - 对关联后的所有
t.TradeAmount做全局去重,涉及大量数据Shuffle(跨节点数据传输)和计算,自然耗时极长。
- 两张表按照
数据量级的放大效应
如果drd.PriceHistoric表中每个PriceDate对应多条记录,Join后的中间数据量会是crm.NewTransaction表的数倍甚至数十倍,后续的去重操作需要处理这些膨胀的数据,进一步拉长了执行时间。而第二个查询因为跳过了Join,只处理原表数据,计算量可以忽略不计。
内容的提问来源于stack exchange,提问作者Patterson
相关产品推荐
相关产品推荐

