Conversation
e510237 to
9169ff0
Compare
Adds a broadcast range join for range-predicate joins that today fall back to BroadcastNestedLoopJoinExec. Point-in-range (for example an IP lookup, `ip BETWEEN low AND high`) broadcasts an interval index over (low, high) and probes it for overlap. Interval overlap (`a.lo < b.hi AND b.lo < a.hi`) uses the same index, with (low, high) on each side, and probes it with overlapping(low, high). A single inequality (for example `a.start < b.end`) broadcasts a sorted point index of that one column. The index has no direction, so `<` and `>` share one broadcast and the join scans the side of the bound that the build side holds. Either index returns a superset. A candidate is emitted only when the original join condition evaluates to true, so inclusivity stays on the predicate. If build-side intervals overlap heavily, the candidate count can approach the build size. Supported types are inner, left outer, right outer, left semi, and left anti. The preserved side is streamed: left outer, left semi, and left anti broadcast the right side, and right outer broadcasts the left side. Full outer stays a broadcast nested loop join. A null streamed key is not a match. The join condition and keys are fields of BroadcastRangeJoinExec, so scalar subqueries in the condition are planned and EXPLAIN prints the original expressions. Keys are bound to the child output at execution. Only types the index compares with the join predicate are recognized; other orderable types stay a broadcast nested loop join. BroadcastRangeJoinExec participates in whole-stage codegen. The stream side stays in the generated stage. Each streamed row boxes its range keys, probes the broadcast index, and evaluates the original condition in generated code. The index search stays interpreted. A non-leaf expression that cannot be generated keeps the interpreted operator. The feature is gated off by default behind spark.sql.join.broadcastRangeJoin.enabled. ExtractRangeJoinKeys recognizes point-in-range, interval overlap, and partial-range shapes, including extra conjuncts. Point-in-range wins over overlap, and overlap wins over a partial range. A deterministic With, which is what BETWEEN lowers to when common expressions are not inlined, is substituted before the conjunct scan so the keys are the original expressions. A non-deterministic definition is not used as a range key. BroadcastRangeJoinExec requires BroadcastDistribution(RangeBroadcastMode), so BroadcastExchangeExec builds the index. Under AQE, LogicalQueryStageStrategy rebuilds the operator when the mode's bound keys match. A BROADCAST hint whose build side the operator supports is planned as a range join. An unsupported hinted side stays a broadcast nested loop join. SHUFFLE_REPLICATE_NL stays a cartesian product. Tested by ExtractRangeJoinKeysSuite (point-in-range, overlap, partial range, and rejected types), RangeIndexSuite (interval sweep, equal-but-not-identical keys, and point-index bounds), RangeJoinSuite (boundary inclusivity, build side, outer, semi, anti, partial upTo/from, and interval overlap), RangeJoinSQLSuite (planner gating, BETWEEN, hints, join types, residual predicates, collation, AQE, and whole-stage codegen), and AdaptiveQueryExecSuite (partition coalescing and range broadcast reuse).
HyukjinKwon
left a comment
There was a problem hiding this comment.
Thanks, this is a nice speedup for a long-standing gap, and keeping the original condition as the only accept/reject check makes the index easy to reason about. One concern about the interval index's broadcast size, inline.
| private[this] val activated: Array[InternalRow]) | ||
| extends RangeRelation { | ||
|
|
||
| override def sizeInBytes(): Long = activated.map(RangeIndex.rowSize).sum |
There was a problem hiding this comment.
I think the persistent-HashMap argument above (line 161) holds only on the driver, before the index is broadcast. Scala 2.13's immutable.HashMap serializes through DefaultSerializationProxy, which writes every entry, so each activeOld/activeAll snapshot is serialized in full and deserialized on each executor as an independent map. For heavily overlapping build intervals (the O(n) keys x O(n) active rows case the comment describes), the serialized broadcast and the executor-side index are then O(n^2), not O(n log n).
sizeInBytes also counts only activated, so spark.sql.maxBroadcastTableSize doesn't see the snapshots. A build side well under the broadcast threshold could still blow up driver serialization or executor memory.
Could sizeInBytes include the snapshot entry counts, or could the snapshots be stored in a form that stays compact on the wire (e.g. an interval tree, or per-key deltas rebuilt lazily on the executor)? A test that serializes an index of n fully overlapping intervals and checks the size would pin this down.
There was a problem hiding this comment.
Thank you Hyukjin. Broadcast only the sweep and rebuild the active-set snapshots on the first probe, so Java serialization no longer expands them to O(n²) now.
There was a problem hiding this comment.
Thanks, rebuilding the snapshots lazily does fix the serialized size, and the round-trip test pins that down. The executor-side half still seems open, though. Each executor now builds the per-key IntMap snapshots on first probe, and they are still O(events x trie depth). I measured about 55 MB for 100k and 244 MB for 400k fully overlapping intervals, roughly 275-305 bytes per sweep event, while sizeInBytes charges 8 bytes per event plus the rows. So spark.sql.maxBroadcastTableSize still doesn't bound what the executor holds. Could sizeInBytes include an estimate for the snapshots (for example, events x log2(active) x node size), or could the probe avoid materializing a snapshot per key (for example, an interval tree or a periodic checkpoint plus a delta scan)?
HyukjinKwon
left a comment
There was a problem hiding this comment.
Review summary
The planner, operator, AQE, and index paths line up well. Keys are gated to types whose index ordering matches the predicate, the original condition stays the only accept/reject check, and the tests compare every supported join type against the nested loop result.
The inline thread about the index's broadcast size is half resolved. Broadcasting only the sweep makes the payload linear, but the snapshots each executor now rebuilds are still invisible to sizeInBytes, and so to spark.sql.maxBroadcastTableSize. I measured roughly 275-305 bytes per sweep event for fully overlapping intervals, against 8 charged. The remaining comments are small doc/code mismatches: UDT keys are documented as supported but rejected, the join-selection scaladoc omits the interval-overlap shape, and one test comment inverts its rationale.
Findings
4 total: 0 P0, 0 P1, 1 P2, 3 P3.
Non-blocking (P2)
- Executor-side interval snapshots are not counted in sizeInBytes —
sql/core/src/main/scala/org/apache/spark/sql/execution/joins/RangeIndex.scala:239— unresolved in existing discussion.
Nit (P3)
- Join-selection scaladoc omits the interval-overlap shape —
sql/core/src/main/scala/org/apache/spark/sql/execution/SparkStrategies.scala:184— see inline. - UDT range keys are documented as supported but rejected —
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/planning/patterns.scala:628— see inline. - Inverted-interval test comment says "rejects" instead of "accepts" —
sql/core/src/test/scala/org/apache/spark/sql/execution/joins/RangeIndexSuite.scala:94— see inline.
Existing discussions
- existing discussion — The review summary flagged the interval index's broadcast size, with details in the inline thread. The current head broadcasts only the sweep, so the serialized payload is now linear. But the lazily rebuilt snapshots still occupy O(events x trie depth) executor heap that sizeInBytes and maxBroadcastTableSize do not see, which is the accounting half of the same concern.
| * Supported for all join types. | ||
| * | ||
| * - Broadcast range join (BRJ): | ||
| * Supports a point-in-range predicate or one cross-side inequality. |
There was a problem hiding this comment.
Nit (P3): Nit: this lists only point-in-range and one cross-side inequality, but ExtractRangeJoinKeys also recognizes interval overlap (a.lo < b.hi AND b.lo < a.hi), and that is planned as BroadcastRangeJoinExec too. Could this mention all three shapes? The empty line a few lines below is also missing its leading * inside the scaladoc block.
Verification:
- Inspection: The BRJ entry lists interval overlap, and no line in that scaladoc block is missing its leading asterisk.
| object RangePredicate { | ||
| /** | ||
| * Types the range index orders with the same comparison as the join predicate. | ||
| * A UDT is ordered as its sql type, which is how the row stores the value. |
There was a problem hiding this comment.
Nit (P3): This says a UDT is ordered as its SQL type, and the PR description lists UDTs as recognized, but supportedType below has no UserDefinedType case, so a UDT key falls through to case _ => false and the join stays a nested loop join. If UDTs should be supported, I think this needs case udt: UserDefinedType[_] => supportedType(udt.sqlType), plus unwrapping the ordering in RangeBroadcastMode.keyOrdering (PhysicalDataType(udt) resolves to UninitializedPhysicalType), and a test. Otherwise could we drop the UDT claim here and in the description?
Verification:
- Behavior: Key extraction recognizes a UDT-over-long predicate, and an end-to-end join over that UDT matches the conf-off result on both build sides.
- Compatibility: Types outside the supported list, including a UDT over an array, remain unrecognized.
|
|
||
| test("an inverted build interval reaches the windows the condition accepts") { | ||
| // (3, 1) is normalized to [1, 3]. A window with low < 1 and high > 3 accepts it, | ||
| // and dropping the row would lose the pair the join condition then rejects. |
There was a problem hiding this comment.
Nit (P3): Nit: I think this should be "accepts". For build (3, 1) and window (0, 4), a.lo < b.hi AND b.lo < a.hi holds, which is why dropping the row would lose a real result (matching the toRangeEvents doc).
Verification:
- Inspection: The test comment matches the toRangeEvents rationale: the pair is accepted by the condition.
What changes were proposed in this pull request?
Add a broadcast range join for range-predicate joins that today fall back to
BroadcastNestedLoopJoinExec. It is off by default (spark.sql.join.broadcastRangeJoin.enabled).ExtractRangeJoinKeysrecognizes three shapes. When more than one matches, point-in-range wins, then interval overlap, then a partial range.Point-in-range. One side is a point and the other is
(low, high), for examplep >= lo AND p <= hiorp BETWEEN lo AND hi. The keys are(point, point)and(low, high). The build side is an interval index, probed withoverlapping(point, point).Interval overlap. Both sides are
(low, high), for examplea.lo < b.hi AND b.lo < a.hi. The keys are(low, high)on each side. The build side is an interval index, probed withoverlapping(low, high).Partial range. One cross-side inequality, one column per side, for example
a < bora > b. The keys are those two columns. The build side is a sorted point index.<and>share that index. The join scansupToorfromaccording to which side holds the lower bound.The index returns a superset. Inclusivity, extra conjuncts, and collation stay on the original condition and are evaluated for each candidate. Null keys are not matches. An equi-join key still wins, so an equality plus a range predicate stays a hash join.
Supported join types are inner, left outer, right outer, left semi, and left anti. The preserved side is streamed. Full outer stays a broadcast nested loop join. A
BROADCASThint is honored only when that build side is supported.SHUFFLE_REPLICATE_NLstays a cartesian product. Types the index can order (including a UDT ordered as its SQL type, and collated strings) are recognized. Arrays, intervals, and other orderable types stay a nested loop join.BroadcastExchangeExecbuilds the index throughRangeBroadcastMode. Under AQE,LogicalQueryStageStrategyrebuilds the operator when that mode matches. Bound-key nullability is normalized so a stage rewrite still satisfiesBroadcastDistribution.This is a broadcast index for a side that already fits the broadcast threshold. It is not a sort-merge range join for two large sides, and there is no
RANGE_JOINhint.Why are the changes needed?
A range predicate such as an IP lookup or an interval overlap has no equality key, so Spark plans
BroadcastNestedLoopJoinExecand compares every streamed row with every build row. A broadcast index probes that same small side in logarithmic time and then rechecks the original predicate.On one of our internal OLAP Spark clusters, we found 7387 SQL statements with a range pattern in one week, 82.95% of the non-equijoin SQL. The largest speedup was 662.4X.
ip BETWEEN begin_ip AND end_ipcal_dt >= start_dt AND cal_dt <= end_dtDoes this PR introduce any user-facing change?
Yes, behind a new configuration that defaults to false, so plans are unchanged unless it is turned on. Compared with released Spark, the configuration does not exist. On master, with it enabled, a supported range predicate is planned as
BroadcastRangeJoinExecinstead ofBroadcastNestedLoopJoinExec.Answers match the nested loop join. Full outer, an equality key, an unsupported type, and a non-deterministic key keep the previous plan. Overlapping build intervals make the index larger than the build table, and the existing
spark.sql.maxBroadcastTableSizecheck still applies to the index.How was this patch tested?
Unit and SQL tests:
ExtractRangeJoinKeysSuite: point-in-range, overlap, partial range, and rejected types.RangeIndexSuite: interval sweep, keys that compare equal but are not identical, and point-index bounds.RangeJoinSuite: boundary inclusivity, both build sides, outer, semi, anti, partialupTo/from, and interval overlap.RangeJoinSQLSuite: the conf off and on,BETWEEN, hints, join types, a residual predicate, collation, AQE, and whole-stage codegen.AdaptiveQueryExecSuite: broadcast reuse, nullability after AQE, build-side validation, and partition coalescing.Local benchmark against broadcast nested loop join (
RangeJoinBenchmark, Apple M2 Max, OpenJDK 17). Selective cases probe 400,000 rows against 4,000 build rows. Best-time ratio:</>, selective build sideUNICODE_CIstringWas this patch authored or co-authored using generative AI tooling?
Generated-by: Cursor (Grok 4.7)