kumarUjjawal commented on code in PR #24805: URL: https://github.com/apache/datafusion/pull/24805#discussion_r3891586869
########## datafusion/sqllogictest/test_files/piecewise_merge_join_matrix.slt: ########## @@ -0,0 +1,676 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at + +# http://www.apache.org/licenses/LICENSE-2.0 + +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# ================================================================== +# PiecewiseMergeJoin correctness across the config matrix +# ================================================================== +# A range join with no equi key is planned as PiecewiseMergeJoin when +# `enable_piecewise_merge_join=true` and as NestedLoopJoin when false. This file +# sweeps that knob {true,false} so every query runs on both, and the two must +# agree row-for-row. It also sweeps `batch_size` {1,2,100,8192} because PWMJ +# sorts the buffered side, scans the streamed side batch by batch, and (for +# existence joins) lowers a shared cross-batch watermark, so the batch layout is +# load-bearing -- batch_size=1 fragments every input into single-row batches, +# 8192 keeps each fixture in one batch. +# +# PWMJ implements the classic joins (Inner/Left/Right/Full) and the LeftSemi / +# LeftAnti existence joins; RightSemi/RightAnti/Mark stay on NestedLoopJoin in +# both combinations. The file has two parts: classic range joins (Part 1) then +# existence joins (Part 2). +# +# Rules for a matrix file: no EXPLAIN (the plan differs per combination), no +# in-file SET of a swept knob, and `rowsort` on every multi-row query so a +# single expected block matches all eight combinations regardless of the order +# each operator emits rows. + +# configMatrix: datafusion.optimizer.enable_piecewise_merge_join=true,false +# configMatrix: datafusion.execution.batch_size=1,2,100,8192 + +# ================================================================== +# Part 1: Classic range joins (Inner / Left / Right / Full) +# ================================================================== +# No equi key and a single range predicate, so each join is a PiecewiseMergeJoin +# when the knob is on and a NestedLoopJoin when off. cj_l.lv = {1,3,5,NULL}, +# cj_r.rv = {3,5,NULL}. A NULL key matches nothing, so it can only appear as an +# unmatched row on an outer side. +statement ok +CREATE TABLE cj_l(lid INT, lv INT); + +statement ok +INSERT INTO cj_l VALUES (1, 1), (2, 3), (3, 5), (4, NULL); + +statement ok +CREATE TABLE cj_r(rid INT, rv INT); + +statement ok +INSERT INTO cj_r VALUES (1, 3), (2, 5), (3, NULL); + +# INNER, all four operators. `<=`/`>=` differ from `<`/`>` through the equal +# values (lv=3=rv, lv=5=rv). +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l JOIN cj_r r ON l.lv < r.rv; +---- +1 1 1 3 +1 1 2 5 +2 3 2 5 + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l JOIN cj_r r ON l.lv <= r.rv; +---- +1 1 1 3 +1 1 2 5 +2 3 1 3 +2 3 2 5 +3 5 2 5 + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l JOIN cj_r r ON l.lv > r.rv; +---- +3 5 1 3 + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l JOIN cj_r r ON l.lv >= r.rv; +---- +2 3 1 3 +3 5 1 3 +3 5 2 5 + +# LEFT: the INNER `<` rows plus every unmatched left row (lv=5 and lv=NULL) with +# a NULL right side. +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l LEFT JOIN cj_r r ON l.lv < r.rv; +---- +1 1 1 3 +1 1 2 5 +2 3 2 5 +3 5 NULL NULL +4 NULL NULL NULL + +# RIGHT: the INNER rows plus every unmatched right row with a NULL left side. The +# unmatched right row rid=3 has a NULL key and must still be emitted (a fixed +# regression dropped streamed-side NULL rows here). +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l RIGHT JOIN cj_r r ON l.lv < r.rv; +---- +1 1 1 3 +1 1 2 5 +2 3 2 5 +NULL NULL 3 NULL + +# RIGHT with a different operator flips the buffered-side sort direction; the +# NULL-keyed right row rid=3 is still emitted unmatched. +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l RIGHT JOIN cj_r r ON l.lv >= r.rv; +---- +2 3 1 3 +3 5 1 3 +3 5 2 5 +NULL NULL 3 NULL + +# FULL: the INNER rows plus unmatched rows from both sides, including a NULL key +# on each side. +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l FULL JOIN cj_r r ON l.lv < r.rv; +---- +1 1 1 3 +1 1 2 5 +2 3 2 5 +3 5 NULL NULL +4 NULL NULL NULL +NULL NULL 3 NULL + +# Empty right side: LEFT keeps every left row with NULLs, INNER is empty. +statement ok +CREATE TABLE cj_empty(rid INT, rv INT); + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l JOIN cj_empty r ON l.lv < r.rv; +---- + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_l l LEFT JOIN cj_empty r ON l.lv < r.rv; +---- +1 1 NULL NULL +2 3 NULL NULL +3 5 NULL NULL +4 NULL NULL NULL + +# Empty left side: RIGHT keeps every right row with NULLs, including rid=3's NULL +# key. +query IIII rowsort +SELECT l.rid, l.rv, r.rid, r.rv FROM cj_empty l RIGHT JOIN cj_r r ON l.rv < r.rv; +---- +NULL NULL 1 3 +NULL NULL 2 5 +NULL NULL 3 NULL + +# Duplicate keys on both sides: each qualifying left row must pair with the whole +# matching right suffix. lv=1 (twice) is below rv=2 (twice) -> 4 rows; lv=3 has +# no larger right value. +statement ok +CREATE TABLE cj_dup_l(lid INT, lv INT); + +statement ok +INSERT INTO cj_dup_l VALUES (1, 1), (2, 1), (3, 3); + +statement ok +CREATE TABLE cj_dup_r(rid INT, rv INT); + +statement ok +INSERT INTO cj_dup_r VALUES (1, 2), (2, 2); + +query IIII rowsort +SELECT l.lid, l.lv, r.rid, r.rv FROM cj_dup_l l JOIN cj_dup_r r ON l.lv < r.rv; +---- +1 1 1 2 +1 1 2 2 +2 1 1 2 +2 1 2 2 + +# ================================================================== +# Part 2: Existence joins (LeftSemi / LeftAnti; Right/Mark fall back to NLJ) +# ================================================================== +# EXISTS with a single range correlation -> LeftSemi, NOT EXISTS -> LeftAnti; +# these two are the only existence joins PWMJ implements. RightSemi/RightAnti +# (explicit syntax) and Mark (`OR EXISTS`) stay on NestedLoopJoin in both +# combinations and are covered at the end. + +# ------------------------------------------------------------------ +# Fixtures +# ------------------------------------------------------------------ +statement ok +CREATE TABLE ej_l(id INT, v INT); + +statement ok +INSERT INTO ej_l VALUES (1, 5), (2, 4), (3, 2), (4, 1); + +statement ok +CREATE TABLE ej_r(v INT); + +statement ok +INSERT INTO ej_r VALUES (2), (3), (4); + +# ------------------------------------------------------------------ +# Basic: every operator, both LeftSemi (EXISTS) and LeftAnti (NOT EXISTS) +# ------------------------------------------------------------------ +# ej_l.v = {5,4,2,1}, ej_r.v = {2,3,4}. + +# `<` : 2<3 and 1<2 match; 5 and 4 exceed every right value. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v < r.v); +---- +3 +4 + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v < r.v); +---- +1 +2 + +# `<=` : additionally admits v=4 through the equal right value 4. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v <= r.v); +---- +2 +3 +4 + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v <= r.v); +---- +1 + +# `>` : 5 and 4 exceed some right value; 2 and 1 do not. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- +1 +2 + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- +3 +4 + +# `>=` : additionally admits v=2 through the equal right value 2. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v >= r.v); +---- +1 +2 +3 + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v >= r.v); +---- +4 + +# ------------------------------------------------------------------ +# Correlation written inner-column-first, and with an expression +# ------------------------------------------------------------------ +# `r.v < l.v` is the same join as `l.v > r.v`; the planner flips the operator so +# the marked (left) side stays buffered. Same rows as `l.v > r.v` above. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE r.v < l.v); +---- +1 +2 + +# An expression on the streamed side: `l.v < r.v + 1` is `l.v <= r.v` over ints. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v < r.v + 1); +---- +2 +3 +4 + +# ------------------------------------------------------------------ +# NULL semantics on the buffered (left) key +# ------------------------------------------------------------------ +# A NULL key satisfies no comparison, so it is excluded from EXISTS and kept by +# NOT EXISTS. ej_ln.v = {NULL,3,NULL,1}, right non-null key = {2}. +statement ok +CREATE TABLE ej_ln(id INT, v INT); + +statement ok +INSERT INTO ej_ln VALUES (1, NULL), (2, 3), (3, NULL), (4, 1); + +statement ok +CREATE TABLE ej_rn(v INT); + +statement ok +INSERT INTO ej_rn VALUES (2), (NULL); + +# `<` : only v=1 is below 2. Both NULL-keyed left rows stay in NOT EXISTS. +query I rowsort +SELECT l.id FROM ej_ln l WHERE EXISTS (SELECT 1 FROM ej_rn r WHERE l.v < r.v); +---- +4 + +query I rowsort +SELECT l.id FROM ej_ln l WHERE NOT EXISTS (SELECT 1 FROM ej_rn r WHERE l.v < r.v); +---- +1 +2 +3 + +# `>` : only v=3 exceeds 2. +query I rowsort +SELECT l.id FROM ej_ln l WHERE EXISTS (SELECT 1 FROM ej_rn r WHERE l.v > r.v); +---- +2 + +query I rowsort +SELECT l.id FROM ej_ln l WHERE NOT EXISTS (SELECT 1 FROM ej_rn r WHERE l.v > r.v); +---- +1 +3 +4 + +# ------------------------------------------------------------------ +# Empty inputs +# ------------------------------------------------------------------ +statement ok +CREATE TABLE ej_empty(v INT); + +# Empty streamed (right) side: no key can match, so EXISTS is empty and NOT +# EXISTS keeps every left row. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_empty r WHERE l.v > r.v); +---- + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_empty r WHERE l.v > r.v); +---- +1 +2 +3 +4 + +statement ok +CREATE TABLE ej_lempty(id INT, v INT); + +# Empty buffered (left) side: nothing to emit either way. +query I rowsort +SELECT l.id FROM ej_lempty l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- + +query I rowsort +SELECT l.id FROM ej_lempty l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- + +# ------------------------------------------------------------------ +# All-NULL sides (distinct early-exit triggers in the existence path) +# ------------------------------------------------------------------ +# An all-NULL buffered side saturates the watermark before the first poll, so no +# streamed batch is read; EXISTS is still empty. +statement ok +CREATE TABLE ej_l_allnull(id INT, v INT); + +statement ok +INSERT INTO ej_l_allnull VALUES (1, NULL), (2, NULL); + +query I rowsort +SELECT l.id FROM ej_l_allnull l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- + +query I rowsort +SELECT l.id FROM ej_l_allnull l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v); +---- +1 +2 + +# An all-NULL streamed side never lowers the watermark, so it behaves like an +# empty streamed side. +statement ok +CREATE TABLE ej_r_allnull(v INT); + +statement ok +INSERT INTO ej_r_allnull VALUES (NULL), (NULL); + +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r_allnull r WHERE l.v > r.v); +---- + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r_allnull r WHERE l.v > r.v); +---- +1 +2 +3 +4 + +# ------------------------------------------------------------------ +# Duplicate buffered keys straddling the match boundary +# ------------------------------------------------------------------ +# The first matching buffered row is found by binary search; it must return the +# FIRST index of a run of equal keys or earlier duplicates vanish from EXISTS. +# ej_dup_l.v = {5,5,3,3,1}; the deciding streamed key is 3, so `<=`/`>=` land +# inside the run of 3s. +statement ok +CREATE TABLE ej_dup_l(id INT, v INT); + +statement ok +INSERT INTO ej_dup_l VALUES (1, 5), (2, 5), (3, 3), (4, 3), (5, 1); + +statement ok +CREATE TABLE ej_dup_r(v INT); + +statement ok +INSERT INTO ej_dup_r VALUES (3); + +# `<` : only v=1 is below 3. +query I rowsort +SELECT l.id FROM ej_dup_l l WHERE EXISTS (SELECT 1 FROM ej_dup_r r WHERE l.v < r.v); +---- +5 + +# `<=` : both v=3 rows must appear. +query I rowsort +SELECT l.id FROM ej_dup_l l WHERE EXISTS (SELECT 1 FROM ej_dup_r r WHERE l.v <= r.v); +---- +3 +4 +5 + +# `>` : only the two v=5 rows exceed 3. +query I rowsort +SELECT l.id FROM ej_dup_l l WHERE EXISTS (SELECT 1 FROM ej_dup_r r WHERE l.v > r.v); +---- +1 +2 + +# `>=` : the boundary is inside the run of 3s from the other direction. +query I rowsort +SELECT l.id FROM ej_dup_l l WHERE EXISTS (SELECT 1 FROM ej_dup_r r WHERE l.v >= r.v); +---- +1 +2 +3 +4 + +query I rowsort +SELECT l.id FROM ej_dup_l l WHERE NOT EXISTS (SELECT 1 FROM ej_dup_r r WHERE l.v >= r.v); +---- +5 + +# ------------------------------------------------------------------ +# Subquery filter -> repartitioned streamed side (multi-partition final pass) +# ------------------------------------------------------------------ +# A predicate on the inner table is pushed to the scan and lets the streamed +# side repartition, so the final existence pass is coordinated across streamed +# partitions. `r.v > 0` keeps all of {2,3,4}, so the rows match `l.v > r.v`. +query I rowsort +SELECT l.id FROM ej_l l WHERE EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v AND r.v > 0); +---- +1 +2 + +query I rowsort +SELECT l.id FROM ej_l l WHERE NOT EXISTS (SELECT 1 FROM ej_r r WHERE l.v > r.v AND r.v > 0); +---- +3 +4 + +# ------------------------------------------------------------------ +# Type coverage: the existence path shares the join comparator +# ------------------------------------------------------------------ + +# Date32 key. +statement ok +CREATE TABLE ej_dl(id INT, d DATE); + +statement ok +INSERT INTO ej_dl VALUES (1, DATE '2022-04-23'), (2, DATE '2022-04-28'), (3, DATE '2022-04-18'); + +statement ok +CREATE TABLE ej_dr(d DATE); + +statement ok +INSERT INTO ej_dr VALUES (DATE '2022-04-20'), (DATE '2022-04-26'); + +query I rowsort +SELECT l.id FROM ej_dl l WHERE EXISTS (SELECT 1 FROM ej_dr r WHERE l.d > r.d); +---- +1 +2 + +query I rowsort +SELECT l.id FROM ej_dl l WHERE NOT EXISTS (SELECT 1 FROM ej_dr r WHERE l.d > r.d); +---- +3 + +# Float key with negative zero: -0.0 and +0.0 compare equal in SQL, so +# `-0.0 < 0.0` is false and `-0.0 <= 0.0` is true. The comparator normalizes the +# sign of zero before comparing. +statement ok +CREATE TABLE ej_fl(id INT, v DOUBLE); + +statement ok +INSERT INTO ej_fl VALUES (1, -0.0), (2, 2.5); + +statement ok +CREATE TABLE ej_fr(v DOUBLE); + +statement ok +INSERT INTO ej_fr VALUES (0.0); + +query I rowsort +SELECT l.id FROM ej_fl l WHERE EXISTS (SELECT 1 FROM ej_fr r WHERE l.v < r.v); +---- + +query I rowsort +SELECT l.id FROM ej_fl l WHERE EXISTS (SELECT 1 FROM ej_fr r WHERE l.v <= r.v); +---- +1 + +# String key. +statement ok +CREATE TABLE ej_sl(id INT, s VARCHAR); + +statement ok +INSERT INTO ej_sl VALUES (1, 'apple'), (2, 'cherry'), (3, 'mango'); + +statement ok +CREATE TABLE ej_sr(s VARCHAR); + +statement ok +INSERT INTO ej_sr VALUES ('banana'), ('lemon'); + +query I rowsort +SELECT l.id FROM ej_sl l WHERE EXISTS (SELECT 1 FROM ej_sr r WHERE l.s > r.s); +---- +2 +3 + +query I rowsort +SELECT l.id FROM ej_sl l WHERE NOT EXISTS (SELECT 1 FROM ej_sr r WHERE l.s > r.s); +---- +1 + +# Dictionary-encoded key: no typed arrow min/max kernel, so the extreme key per +# streamed batch is chosen through the generic ScalarValue path. +statement ok +CREATE TABLE ej_dict_l AS + SELECT column1 AS id, arrow_cast(column2, 'Dictionary(Int32, Utf8)') AS v + FROM (VALUES (1, 'a'), (2, 'c'), (3, 'e'), (4, NULL)); + +statement ok +CREATE TABLE ej_dict_r AS + SELECT arrow_cast(column1, 'Dictionary(Int32, Utf8)') AS v FROM (VALUES ('c')); + +query I rowsort +SELECT l.id FROM ej_dict_l l WHERE EXISTS (SELECT 1 FROM ej_dict_r r WHERE l.v > r.v); +---- +3 + +query I rowsort +SELECT l.id FROM ej_dict_l l WHERE NOT EXISTS (SELECT 1 FROM ej_dict_r r WHERE l.v > r.v); +---- +1 +2 +4 + +query I rowsort +SELECT l.id FROM ej_dict_l l WHERE EXISTS (SELECT 1 FROM ej_dict_r r WHERE l.v <= r.v); +---- +1 +2 + +# ------------------------------------------------------------------ +# Multi-batch stress (verified by count) +# ------------------------------------------------------------------ +# Larger inputs so batch_size=1 and 2 split each side into many batches, driving Review Comment: batch_size=1 does not split this VALUES input into one-row batches. This matrix does not cover the stated per-batch extreme-key path or the cross-batch watermark path. Can we build the streamed fixture from a source that emits several batches, such as generate_series and add a plan or metric assertion. -- 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]
