Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -700,10 +700,12 @@ jobs:
rust-version: stable
- name: Run sqllogictest
# TODO: Right now several tests are failing in Substrait round-trip mode, so this
# command cannot be run for all the .slt files. Run it for just one that works (limit.slt)
# command cannot be run for all the .slt files. Run a supported subset
# until most of the tickets in https://github.com/apache/datafusion/issues/16248 are addressed
# and this command can be run without filters.
run: cargo xtask ci step test substrait
run: |
cargo xtask ci step test substrait
cargo xtask ci step test substrait-optimized

# Temporarily commenting out the Windows flow, the reason is enormously slow running build
# Waiting for new Windows 2025 github runner
Expand Down
11 changes: 10 additions & 1 deletion datafusion/sqllogictest/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -437,7 +437,8 @@ Not all statements will be round-tripped, some statements like CREATE, INSERT, S
issued as is, but any other statement will be round-tripped to/from Substrait.

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

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

Add `--substrait-optimize` to optimize the logical plan before serialization. This exercises
producer inputs created by optimizer rewrites, such as conditionless joins from `LEFT JOIN ... ON true`
and uncorrelated `WHERE EXISTS`:

```shell
cargo test --test sqllogictests --features substrait -- --substrait-round-trip --substrait-optimize joins_conditionless.slt
```

## `.slt` file format

[`sqllogictest`] was originally written for SQLite to verify the
Expand Down
29 changes: 24 additions & 5 deletions datafusion/sqllogictest/bin/sqllogictests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,7 @@ async fn run_tests() -> Result<()> {
filters.as_ref(),
currently_running_sql_tracker_clone,
colored_output,
options.substrait_optimize,
)
.await
}
Expand Down Expand Up @@ -461,7 +462,9 @@ fn is_env_truthy(name: &str) -> bool {
enum Engine {
DataFusion,
#[cfg(feature = "substrait")]
SubstraitRoundTrip,
SubstraitRoundTrip {
optimize: bool,
},
}

impl Engine {
Expand All @@ -471,12 +474,13 @@ impl Engine {
match self {
Engine::DataFusion => "Datafusion",
#[cfg(feature = "substrait")]
Engine::SubstraitRoundTrip => "DatafusionSubstraitRoundTrip",
Engine::SubstraitRoundTrip { .. } => "DatafusionSubstraitRoundTrip",
}
}
}

#[cfg(feature = "substrait")]
#[expect(clippy::too_many_arguments, reason = "mirrors the other file runners")]
async fn run_test_file_substrait_round_trip(
test_file: TestFile,
validator: Validator,
Expand All @@ -485,9 +489,10 @@ async fn run_test_file_substrait_round_trip(
filters: &[Filter],
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
colored_output: bool,
optimize: bool,
) -> Result<()> {
run_matrix(
Engine::SubstraitRoundTrip,
Engine::SubstraitRoundTrip { optimize },
test_file,
validator,
mp,
Expand All @@ -504,6 +509,10 @@ async fn run_test_file_substrait_round_trip(
clippy::unused_async,
reason = "matches the substrait-enabled implementation"
)]
#[expect(
clippy::too_many_arguments,
reason = "matches the enabled implementation"
)]
async fn run_test_file_substrait_round_trip(
_test_file: TestFile,
_validator: Validator,
Expand All @@ -512,6 +521,7 @@ async fn run_test_file_substrait_round_trip(
_filters: &[Filter],
_currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
_colored_output: bool,
_optimize: bool,
) -> Result<()> {
exec_err!("Cannot run substrait round-trip: the 'substrait' feature is not enabled")
}
Expand Down Expand Up @@ -620,8 +630,8 @@ impl MatrixRunner<'_> {
match self.engine {
Engine::DataFusion => self.run_datafusion(&test_ctx, pb).await,
#[cfg(feature = "substrait")]
Engine::SubstraitRoundTrip => {
self.run_substrait_round_trip(&test_ctx, pb).await
Engine::SubstraitRoundTrip { optimize } => {
self.run_substrait_round_trip(&test_ctx, pb, optimize).await
}
}
}
Expand Down Expand Up @@ -689,13 +699,15 @@ impl MatrixRunner<'_> {
&self,
test_ctx: &TestContext,
pb: ProgressBar,
optimize: bool,
) -> Result<()> {
let mut runner = sqllogictest::Runner::new(|| async {
Ok(DataFusionSubstraitRoundTrip::new(
test_ctx.session_ctx().clone(),
self.relative_path.to_path_buf(),
pb.clone(),
)
.with_optimization(optimize)
.with_currently_executing_sql_tracker(
self.currently_executing_sql_tracker.clone(),
))
Expand Down Expand Up @@ -1065,6 +1077,13 @@ struct Options {
)]
substrait_round_trip: bool,

#[clap(
long,
requires = "substrait_round_trip",
help = "Optimize logical plans before serializing them in Substrait round-trip mode"
)]
substrait_optimize: bool,

#[clap(long, env = "INCLUDE_SQLITE", help = "Include sqlite files")]
include_sqlite: bool,

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ pub struct DataFusionSubstraitRoundTrip {
relative_path: PathBuf,
pb: ProgressBar,
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker,
optimize_before_round_trip: bool,
}

impl DataFusionSubstraitRoundTrip {
Expand All @@ -51,9 +52,16 @@ impl DataFusionSubstraitRoundTrip {
relative_path,
pb,
currently_executing_sql_tracker: CurrentlyExecutingSqlTracker::default(),
optimize_before_round_trip: false,
}
}

/// Optimize the logical plan before serializing it to Substrait.
pub fn with_optimization(mut self, optimize: bool) -> Self {
self.optimize_before_round_trip = optimize;
self
}

/// Add a tracker that will track the currently executed SQL statement.
///
/// This is useful for logging and debugging purposes.
Expand Down Expand Up @@ -101,7 +109,12 @@ impl sqllogictest::AsyncDB for DataFusionSubstraitRoundTrip {
let tracked_sql = self.currently_executing_sql_tracker.set_sql(sql);

let start = Instant::now();
let result = run_query_substrait_round_trip(&self.ctx, sql).await;
let result = run_query_substrait_round_trip(
&self.ctx,
sql,
self.optimize_before_round_trip,
)
.await;
let duration = start.elapsed();
rewrap_replaced_pool(&self.ctx, &self.relative_path.display().to_string());

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

let state = ctx.state();
let round_tripped_plan = match df.logical_plan() {
let logical_plan = if optimize {
df.into_optimized_plan()?
} else {
df.into_unoptimized_plan()
};
let round_tripped_plan = match &logical_plan {
// Substrait does not handle these plans
LogicalPlan::Ddl(_)
| LogicalPlan::Explain(_)
| LogicalPlan::Dml(_)
| LogicalPlan::Copy(_)
| LogicalPlan::DescribeTable(_)
| LogicalPlan::Statement(_) => df.logical_plan().clone(),
| LogicalPlan::Statement(_) => logical_plan,
// For any other plan, convert to Substrait
logical_plan => {
let plan = to_substrait_plan(logical_plan, &state)?;
Expand Down
77 changes: 77 additions & 0 deletions datafusion/sqllogictest/test_files/joins_conditionless.slt
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
# 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.

statement ok
CREATE TABLE conditionless_left(a BIGINT) AS VALUES (1), (2);

statement ok
CREATE TABLE conditionless_right(b BIGINT) AS VALUES (10), (20);

statement ok
CREATE TABLE conditionless_empty(b BIGINT);

# Optimization removes the true filter. Substrait still requires a join condition.
query II rowsort
SELECT a, b FROM conditionless_left LEFT JOIN conditionless_right ON true;
----
1 10
1 20
2 10
2 20

query II rowsort
SELECT a, b FROM conditionless_left RIGHT JOIN conditionless_right ON true;
----
1 10
1 20
2 10
2 20

query II rowsort
SELECT a, b FROM conditionless_left FULL JOIN conditionless_right ON true;
----
1 10
1 20
2 10
2 20

query II rowsort
SELECT a, b FROM conditionless_left LEFT JOIN conditionless_empty ON true;
----
1 NULL
2 NULL

# Uncorrelated EXISTS and NOT EXISTS become conditionless semi and anti joins.
query I rowsort
SELECT a FROM conditionless_left WHERE EXISTS (SELECT 1 FROM conditionless_right);
----
1
2

query I rowsort
SELECT a FROM conditionless_left WHERE EXISTS (SELECT 1 FROM conditionless_empty);
----

query I rowsort
SELECT a FROM conditionless_left WHERE NOT EXISTS (SELECT 1 FROM conditionless_right);
----

query I rowsort
SELECT a FROM conditionless_left WHERE NOT EXISTS (SELECT 1 FROM conditionless_empty);
----
1
2
Loading
Loading