Skip to content

Commit f9b7f34

Browse files
fix(substrait): roundtrip conditionless joins (#23469)
## Which issue does this PR close? - Closes #18066. ## Rationale for this change Queries with conditionless joins can fail when their optimized logical plans roundtrip through Substrait, with `Plan("join condition should not be empty")`. This affects `LEFT JOIN ... ON true` and uncorrelated `WHERE EXISTS`, as well as scalar-subquery projections when the legacy scalar-subquery rewrite is enabled. The optimizer can remove a constant `true` join filter. The producer then omitted `JoinRel.expression`, which Substrait requires. ## What changes are included in this PR? - Serialize conditionless inner joins as `CrossRel` and other conditionless joins with a literal `true` condition. - Handle overlapping input names when producing the synthetic condition for semi, anti, and mark joins. - Add SQL result coverage for outer joins and uncorrelated `EXISTS` / `NOT EXISTS`, including empty inputs. - Add an opt-in `--substrait-optimize` test mode and run the new fixture through it in CI. This exercises the optimized plans that trigger the bug. ## What is the testing strategy for this PR? `joins_conditionless.slt` checks query results through ordinary SQL execution and optimized Substrait roundtripping. Existing producer and DataFrame tests cover conditionless inner joins, scalar-subquery projections, and overlapping input names. Local validation passed: formatting, Clippy with all targets and features, the extended workspace suite, and `./dev/rust_lint.sh`. The extended suite ran in an isolated PID namespace so the RSS tests could sample a valid baseline on this host. Removing the producer fix makes seven of the new optimized roundtrip cases fail with `join condition should not be empty`; all pass with the fix restored. ## Are there any user-facing changes? Queries containing conditionless joins can roundtrip through Substrait successfully. No breaking production API changes. --------- Co-authored-by: Bruno Volpato <bruno.volpato@datadoghq.com>
1 parent 1fa378d commit f9b7f34

8 files changed

Lines changed: 320 additions & 26 deletions

File tree

‎.github/workflows/rust.yml‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -700,10 +700,12 @@ jobs:
700700
rust-version: stable
701701
- name: Run sqllogictest
702702
# TODO: Right now several tests are failing in Substrait round-trip mode, so this
703-
# command cannot be run for all the .slt files. Run it for just one that works (limit.slt)
703+
# command cannot be run for all the .slt files. Run a supported subset
704704
# until most of the tickets in https://github.com/apache/datafusion/issues/16248 are addressed
705705
# and this command can be run without filters.
706-
run: cargo xtask ci step test substrait
706+
run: |
707+
cargo xtask ci step test substrait
708+
cargo xtask ci step test substrait-optimized
707709
708710
# Temporarily commenting out the Windows flow, the reason is enormously slow running build
709711
# Waiting for new Windows 2025 github runner

‎datafusion/sqllogictest/README.md‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -437,7 +437,8 @@ Not all statements will be round-tripped, some statements like CREATE, INSERT, S
437437
issued as is, but any other statement will be round-tripped to/from Substrait.
438438

439439
_WARNING_: this mode lives behind the `substrait` feature, and the full suite still reports failures. CI therefore
440-
runs it over a single file, through `cargo xtask ci step test substrait`, which filters to `limit.slt`. Some of the
440+
runs `limit.slt` through `cargo xtask ci step test substrait` and optimized conditionless-join tests through
441+
`cargo xtask ci step test substrait-optimized`. Some of the
441442
failures are collected in https://github.com/apache/datafusion/issues/16248. To run the default suite in this mode:
442443

443444
```shell
@@ -450,6 +451,14 @@ For focusing on one specific failing test, a file:line filter can be used:
450451
cargo test --test sqllogictests --features substrait -- --substrait-round-trip binary.slt:23
451452
```
452453

454+
Add `--substrait-optimize` to optimize the logical plan before serialization. This exercises
455+
producer inputs created by optimizer rewrites, such as conditionless joins from `LEFT JOIN ... ON true`
456+
and uncorrelated `WHERE EXISTS`:
457+
458+
```shell
459+
cargo test --test sqllogictests --features substrait -- --substrait-round-trip --substrait-optimize joins_conditionless.slt
460+
```
461+
453462
## `.slt` file format
454463

455464
[`sqllogictest`] was originally written for SQLite to verify the

‎datafusion/sqllogictest/bin/sqllogictests.rs‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -239,6 +239,7 @@ async fn run_tests() -> Result<()> {
239239
filters.as_ref(),
240240
currently_running_sql_tracker_clone,
241241
colored_output,
242+
options.substrait_optimize,
242243
)
243244
.await
244245
}
@@ -461,7 +462,9 @@ fn is_env_truthy(name: &str) -> bool {
461462
enum Engine {
462463
DataFusion,
463464
#[cfg(feature = "substrait")]
464-
SubstraitRoundTrip,
465+
SubstraitRoundTrip {
466+
optimize: bool,
467+
},
465468
}
466469

467470
impl Engine {
@@ -471,12 +474,13 @@ impl Engine {
471474
match self {
472475
Engine::DataFusion => "Datafusion",
473476
#[cfg(feature = "substrait")]
474-
Engine::SubstraitRoundTrip => "DatafusionSubstraitRoundTrip",
477+
Engine::SubstraitRoundTrip { .. } => "DatafusionSubstraitRoundTrip",
475478
}
476479
}
477480
}
478481

479482
#[cfg(feature = "substrait")]
483+
#[expect(clippy::too_many_arguments, reason = "mirrors the other file runners")]
480484
async fn run_test_file_substrait_round_trip(
481485
test_file: TestFile,
482486
validator: Validator,
@@ -485,9 +489,10 @@ async fn run_test_file_substrait_round_trip(
485489
filters: &[Filter],
486490
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
487491
colored_output: bool,
492+
optimize: bool,
488493
) -> Result<()> {
489494
run_matrix(
490-
Engine::SubstraitRoundTrip,
495+
Engine::SubstraitRoundTrip { optimize },
491496
test_file,
492497
validator,
493498
mp,
@@ -504,6 +509,10 @@ async fn run_test_file_substrait_round_trip(
504509
clippy::unused_async,
505510
reason = "matches the substrait-enabled implementation"
506511
)]
512+
#[expect(
513+
clippy::too_many_arguments,
514+
reason = "matches the enabled implementation"
515+
)]
507516
async fn run_test_file_substrait_round_trip(
508517
_test_file: TestFile,
509518
_validator: Validator,
@@ -512,6 +521,7 @@ async fn run_test_file_substrait_round_trip(
512521
_filters: &[Filter],
513522
_currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
514523
_colored_output: bool,
524+
_optimize: bool,
515525
) -> Result<()> {
516526
exec_err!("Cannot run substrait round-trip: the 'substrait' feature is not enabled")
517527
}
@@ -620,8 +630,8 @@ impl MatrixRunner<'_> {
620630
match self.engine {
621631
Engine::DataFusion => self.run_datafusion(&test_ctx, pb).await,
622632
#[cfg(feature = "substrait")]
623-
Engine::SubstraitRoundTrip => {
624-
self.run_substrait_round_trip(&test_ctx, pb).await
633+
Engine::SubstraitRoundTrip { optimize } => {
634+
self.run_substrait_round_trip(&test_ctx, pb, optimize).await
625635
}
626636
}
627637
}
@@ -689,13 +699,15 @@ impl MatrixRunner<'_> {
689699
&self,
690700
test_ctx: &TestContext,
691701
pb: ProgressBar,
702+
optimize: bool,
692703
) -> Result<()> {
693704
let mut runner = sqllogictest::Runner::new(|| async {
694705
Ok(DataFusionSubstraitRoundTrip::new(
695706
test_ctx.session_ctx().clone(),
696707
self.relative_path.to_path_buf(),
697708
pb.clone(),
698709
)
710+
.with_optimization(optimize)
699711
.with_currently_executing_sql_tracker(
700712
self.currently_executing_sql_tracker.clone(),
701713
))
@@ -1065,6 +1077,13 @@ struct Options {
10651077
)]
10661078
substrait_round_trip: bool,
10671079

1080+
#[clap(
1081+
long,
1082+
requires = "substrait_round_trip",
1083+
help = "Optimize logical plans before serializing them in Substrait round-trip mode"
1084+
)]
1085+
substrait_optimize: bool,
1086+
10681087
#[clap(long, env = "INCLUDE_SQLITE", help = "Include sqlite files")]
10691088
include_sqlite: bool,
10701089

‎datafusion/sqllogictest/src/engines/datafusion_substrait_roundtrip_engine/runner.rs‎

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ pub struct DataFusionSubstraitRoundTrip {
4242
relative_path: PathBuf,
4343
pb: ProgressBar,
4444
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
45+
optimize_before_round_trip: bool,
4546
}
4647

4748
impl DataFusionSubstraitRoundTrip {
@@ -51,9 +52,16 @@ impl DataFusionSubstraitRoundTrip {
5152
relative_path,
5253
pb,
5354
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker::default(),
55+
optimize_before_round_trip: false,
5456
}
5557
}
5658

59+
/// Optimize the logical plan before serializing it to Substrait.
60+
pub fn with_optimization(mut self, optimize: bool) -> Self {
61+
self.optimize_before_round_trip = optimize;
62+
self
63+
}
64+
5765
/// Add a tracker that will track the currently executed SQL statement.
5866
///
5967
/// This is useful for logging and debugging purposes.
@@ -101,7 +109,12 @@ impl sqllogictest::AsyncDB for DataFusionSubstraitRoundTrip {
101109
let tracked_sql = self.currently_executing_sql_tracker.set_sql(sql);
102110

103111
let start = Instant::now();
104-
let result = run_query_substrait_round_trip(&self.ctx, sql).await;
112+
let result = run_query_substrait_round_trip(
113+
&self.ctx,
114+
sql,
115+
self.optimize_before_round_trip,
116+
)
117+
.await;
105118
let duration = start.elapsed();
106119
rewrap_replaced_pool(&self.ctx, &self.relative_path.display().to_string());
107120

@@ -143,19 +156,25 @@ impl sqllogictest::AsyncDB for DataFusionSubstraitRoundTrip {
143156
async fn run_query_substrait_round_trip(
144157
ctx: &SessionContext,
145158
sql: impl Into<String>,
159+
optimize: bool,
146160
) -> Result<DFOutput> {
147161
let df = ctx.sql(sql.into().as_str()).await?;
148162
let task_ctx = Arc::new(df.task_ctx());
149163

150164
let state = ctx.state();
151-
let round_tripped_plan = match df.logical_plan() {
165+
let logical_plan = if optimize {
166+
df.into_optimized_plan()?
167+
} else {
168+
df.into_unoptimized_plan()
169+
};
170+
let round_tripped_plan = match &logical_plan {
152171
// Substrait does not handle these plans
153172
LogicalPlan::Ddl(_)
154173
| LogicalPlan::Explain(_)
155174
| LogicalPlan::Dml(_)
156175
| LogicalPlan::Copy(_)
157176
| LogicalPlan::DescribeTable(_)
158-
| LogicalPlan::Statement(_) => df.logical_plan().clone(),
177+
| LogicalPlan::Statement(_) => logical_plan,
159178
// For any other plan, convert to Substrait
160179
logical_plan => {
161180
let plan = to_substrait_plan(logical_plan, &state)?;
Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
# Licensed to the Apache Software Foundation (ASF) under one
2+
# or more contributor license agreements. See the NOTICE file
3+
# distributed with this work for additional information
4+
# regarding copyright ownership. The ASF licenses this file
5+
# to you under the Apache License, Version 2.0 (the
6+
# "License"); you may not use this file except in compliance
7+
# with the License. You may obtain a copy of the License at
8+
#
9+
# http://www.apache.org/licenses/LICENSE-2.0
10+
#
11+
# Unless required by applicable law or agreed to in writing,
12+
# software distributed under the License is distributed on an
13+
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
# KIND, either express or implied. See the License for the
15+
# specific language governing permissions and limitations
16+
# under the License.
17+
18+
statement ok
19+
CREATE TABLE conditionless_left(a BIGINT) AS VALUES (1), (2);
20+
21+
statement ok
22+
CREATE TABLE conditionless_right(b BIGINT) AS VALUES (10), (20);
23+
24+
statement ok
25+
CREATE TABLE conditionless_empty(b BIGINT);
26+
27+
# Optimization removes the true filter. Substrait still requires a join condition.
28+
query II rowsort
29+
SELECT a, b FROM conditionless_left LEFT JOIN conditionless_right ON true;
30+
----
31+
1 10
32+
1 20
33+
2 10
34+
2 20
35+
36+
query II rowsort
37+
SELECT a, b FROM conditionless_left RIGHT JOIN conditionless_right ON true;
38+
----
39+
1 10
40+
1 20
41+
2 10
42+
2 20
43+
44+
query II rowsort
45+
SELECT a, b FROM conditionless_left FULL JOIN conditionless_right ON true;
46+
----
47+
1 10
48+
1 20
49+
2 10
50+
2 20
51+
52+
query II rowsort
53+
SELECT a, b FROM conditionless_left LEFT JOIN conditionless_empty ON true;
54+
----
55+
1 NULL
56+
2 NULL
57+
58+
# Uncorrelated EXISTS and NOT EXISTS become conditionless semi and anti joins.
59+
query I rowsort
60+
SELECT a FROM conditionless_left WHERE EXISTS (SELECT 1 FROM conditionless_right);
61+
----
62+
1
63+
2
64+
65+
query I rowsort
66+
SELECT a FROM conditionless_left WHERE EXISTS (SELECT 1 FROM conditionless_empty);
67+
----
68+
69+
query I rowsort
70+
SELECT a FROM conditionless_left WHERE NOT EXISTS (SELECT 1 FROM conditionless_right);
71+
----
72+
73+
query I rowsort
74+
SELECT a FROM conditionless_left WHERE NOT EXISTS (SELECT 1 FROM conditionless_empty);
75+
----
76+
1
77+
2

0 commit comments

Comments
 (0)