Spark SQL中过滤条件位置对千万级表查询性能的影响咨询
Great question! Let's dive into how Spark SQL handles these two approaches, especially given your specific setup where date_dim only holds a single row of daily data.
Key Context: Spark's Catalyst Optimizer
First, it's critical to remember that Spark uses the Catalyst Optimizer to rewrite and optimize your query's logical and physical execution plans. A core optimization here is predicate pushdown—pushing filtering conditions as early as possible in the execution pipeline to reduce the amount of data processed in subsequent steps (like joins).
Analysis of Your Two Schemes
Scheme 1: Filter in JOIN ON Clause with Subquery
select a.id, b.name, c.salary from tableA a inner join tableB b on a.id = b.id and a.eff_dt <= (select last_mth_day from date_dim) inner join tableC c on b.name = c.name
- Since
date_dimonly has one row, the subquery(select last_mth_day from date_dim)is a scalar subquery. Spark will evaluate this once upfront to get a constant value (e.g.,'31/05/2019'). - The Catalyst optimizer will recognize that
a.eff_dt <= [constant]can be pushed down directly to the scan oftableA. This means Spark will first filtertableAto only rows whereeff_dtmeets the condition, then perform the join withtableB—reducing the dataset size early.
Scheme 2: Filter in WHERE Clause with Cross Join
select a.id, b.name, c.salary from tableA a inner join tableB b on a.id = b.id inner join tableC c on b.name = c.name cross join date_dim dt where a.eff_dt <= dt.last_mth_day
- At first glance, it looks like you're joining all three tables first, then cross-joining with
date_dimand filtering. But becausedate_dimhas only one row, the cross join doesn't inflate your dataset (it just adds a single column with the same value to every row). - Even better, Catalyst will rewrite this query: it will extract the
a.eff_dt <= dt.last_mth_daycondition, resolvedt.last_mth_dayto a constant (thanks to the single-rowdate_dim), and push this filter down to thetableAscan—just like in Scheme 1. The cross join will effectively be optimized away, since it doesn't change the row count or require additional processing.
Performance Verdict
In your specific scenario:
- There will be no significant performance difference between the two schemes. The Catalyst optimizer will normalize both queries into nearly identical physical execution plans, with the filtering happening at the earliest possible stage (scanning
tableA). - If you want to confirm this, run
EXPLAIN EXTENDEDon both queries—you'll see that the optimized logical and physical plans are almost identical.
Readability & Maintainability
While performance is the same, Scheme 1 is more readable. It makes it explicit that the filter is tied to the join between tableA and tableB, whereas Scheme 2's cross join might confuse future maintainers (who might wonder why you're joining with date_dim at all, unless they notice it only has one row).
Edge Case to Note
If date_dim ever had multiple rows, Scheme 2 would be problematic: the cross join would multiply your dataset size by the number of rows in date_dim, leading to unnecessary data processing. But since you've confirmed date_dim only holds one row daily, this isn't a concern here.
内容的提问来源于stack exchange,提问作者marie20

