1fanwang opened a new pull request, #25767:
URL: https://github.com/apache/datafusion/pull/25767

   ## Which issue does this PR close?
   
   - Related to #12955 (plan shape question).
   - Stacked on #25742, which should merge first. Changes after its head 
84ef8f7 are this PR.
   
   ## Rationale for this change
   
   With #25742, `INTERSECT ALL` and `EXCEPT ALL` are about 20x slower than 
`INTERSECT` on mostly distinct rows, because a `row_number()` window makes each 
row its own partition. kita-renji measured the same gap in that review.
   
   ## What changes are included in this PR?
   
   Each side counts its distinct rows with a hash aggregate. The counts are 
joined NULL-safely on the row values, and each row is repeated `min(m, n)` or 
`m - n` times with `unnest(range(..))`. Grouping treats 0.0 and -0.0 as equal, 
which fixes the #25766 counts here.
   
   `LogicalPlanBuilder::intersect_all` and `except_all` (unreleased) now take 
`count` and `range` instead of `row_number`.
   
   The Substrait producer cannot serialize `Unnest`, so SQL `INTERSECT ALL` and 
`EXCEPT ALL` can no longer be produced as Substrait. Consuming still works. 
This is the main cost.
   
   ## What is the testing strategy for this PR?
   
   `datafusion-cli -f bench.sql`, release-nonlto, 18 cores:
   
   ```sql
   CREATE TABLE l AS SELECT value AS a, concat('v', value) AS b FROM 
generate_series(1, 10000000);
   CREATE TABLE r AS SELECT value AS a, concat('v', value) AS b FROM 
generate_series(5000001, 15000000);
   CREATE TABLE ld AS SELECT value % 1000 AS a, concat('v', value % 1000) AS b 
FROM generate_series(1, 10000000);
   CREATE TABLE rd AS SELECT value % 1000 AS a, concat('v', value % 1000) AS b 
FROM generate_series(5000001, 15000000);
   SELECT count(*) FROM (SELECT a, b FROM l INTERSECT ALL SELECT a, b FROM r) t;
   SELECT count(*) FROM (SELECT a, b FROM l EXCEPT ALL SELECT a, b FROM r) t;
   SELECT count(*) FROM (SELECT a, b FROM l INTERSECT SELECT a, b FROM r) t;
   SELECT count(*) FROM (SELECT a, b FROM ld INTERSECT ALL SELECT a, b FROM rd) 
t;
   SELECT count(*) FROM (SELECT a, b FROM ld EXCEPT ALL SELECT a, b FROM rd) t;
   EXPLAIN ANALYZE SELECT count(*) FROM (SELECT a, b FROM l INTERSECT ALL 
SELECT a, b FROM r) t;
   ```
   
   ```
   query                                   #25742 (84ef8f7)     this PR
   l INTERSECT ALL r                       2.264s, 2.313s       0.157s, 0.122s
   l EXCEPT ALL r                          2.496s, 2.387s       0.122s, 0.123s
   l INTERSECT r                           0.085s               0.079s
   ld INTERSECT ALL rd                     0.204s               0.021s
   ld EXCEPT ALL rd                        0.214s               0.017s
   ```
   
   Row counts match. `EXPLAIN ANALYZE` on 84ef8f7:
   
   ```
   BoundedWindowAggExec: ..., mode=[Sorted], metrics=[output_rows=10.00 M, 
elapsed_compute=20.32s, ...]
   BoundedWindowAggExec: ..., mode=[Sorted], metrics=[output_rows=10.00 M, 
elapsed_compute=20.17s, ...]
   ```
   
   With this PR:
   
   ```
   UnnestExec, metrics=[output_rows=5.00 M, elapsed_compute=37.73ms, ...]
     ProjectionExec: expr=[range(CASE WHEN 
__datafusion_set_operation_right_count@0 > 
__datafusion_set_operation_left_count@1 THEN 
__datafusion_set_operation_left_count@1 ELSE 
__datafusion_set_operation_right_count@0 END) as 
__datafusion_set_operation_copies], ...
       HashJoinExec: mode=Partitioned, join_type=Inner, on=[(a@0, a@0), (b@1, 
b@1)], NullsEqual: true, metrics=[output_rows=5.00 M, elapsed_compute=909.84ms, 
...]
         AggregateExec: mode=FinalPartitioned, gby=[a@0 as a, b@1 as b], 
aggr=[count(1) as __datafusion_set_operation_left_count], 
metrics=[output_rows=10.00 M, elapsed_compute=439.33ms, ...]
         AggregateExec: mode=FinalPartitioned, gby=[a@0 as a, b@1 as b], 
aggr=[count(1) as __datafusion_set_operation_right_count], 
metrics=[output_rows=10.00 M, elapsed_compute=453.54ms, ...]
   ```
   
   Signed zero:
   
   ```sql
   SELECT a FROM (VALUES (0.0), (-0.0)) t(a) INTERSECT ALL SELECT a FROM 
(VALUES (0.0)) s(a);
   -- 84ef8f7: -0.0, 0.0    this PR: 0.0
   SELECT a FROM (VALUES (0.0), (-0.0)) t(a) EXCEPT ALL SELECT a FROM (VALUES 
(0.0)) s(a);
   -- 84ef8f7: no rows      this PR: 0.0
   ```
   
   `negative_zero.slt` covers both.
   
   The `self_referential_*_all` Substrait tests now assert `Unsupported plan 
type: Unnest`.
   
   ## Are there any user-facing changes?
   
   Faster `INTERSECT ALL` and `EXCEPT ALL`. They need `count` and `range` 
registered (upgrade guide updated) and can't be produced as Substrait.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to