diff --git a/datafusion/sql/src/query.rs b/datafusion/sql/src/query.rs index e2b9e4d2d5305..0fcf3fb75c19c 100644 --- a/datafusion/sql/src/query.rs +++ b/datafusion/sql/src/query.rs @@ -24,12 +24,12 @@ use datafusion_expr::expr::{Sort, WildcardOptions}; use datafusion_expr::select_expr::SelectExpr; use datafusion_expr::{ - CreateMemoryTable, DdlStatement, Distinct, Expr, LogicalPlan, LogicalPlanBuilder, + CreateMemoryTable, DdlStatement, Distinct, Expr, LogicalPlan, LogicalPlanBuilder, lit, }; use sqlparser::ast::{ - Expr as SQLExpr, ExprWithAliasAndOrderBy, Ident, LimitClause, Offset, OffsetRows, - OrderBy, OrderByExpr, OrderByKind, PipeOperator, Query, SelectInto, SetExpr, - SetOperator, SetQuantifier, TableAlias, + Expr as SQLExpr, ExprWithAliasAndOrderBy, Fetch, Ident, LimitClause, Offset, + OffsetRows, OrderBy, OrderByExpr, OrderByKind, PipeOperator, Query, SelectInto, + SetExpr, SetOperator, SetQuantifier, TableAlias, }; use sqlparser::tokenizer::Span; @@ -58,10 +58,6 @@ impl SqlToRel<'_, S> { pipe_operators, } = query; - if fetch.is_some() { - return not_impl_err!("FETCH clause is not supported yet"); - } - if let Some(with) = with { self.plan_with_clause(with, planner_context)?; } @@ -72,7 +68,12 @@ impl SqlToRel<'_, S> { let select_into = select.into.take(); let plan = self.select_to_plan(*select, order_by.clone(), planner_context)?; - let plan = self.limit(plan, limit_clause.clone(), planner_context)?; + let plan = self.limit( + plan, + limit_clause.clone(), + fetch.clone(), + planner_context, + )?; // Process the `SELECT INTO` after `LIMIT`. self.select_into(plan, select_into) } @@ -92,7 +93,7 @@ impl SqlToRel<'_, S> { None, )?; let plan = self.order_by(plan, order_by_rex)?; - self.limit(plan, limit_clause, planner_context) + self.limit(plan, limit_clause, fetch, planner_context) } }?; @@ -143,6 +144,7 @@ impl SqlToRel<'_, S> { }), limit_by: vec![], }), + None, planner_context, ), PipeOperator::Select { exprs } => { @@ -245,20 +247,21 @@ impl SqlToRel<'_, S> { &self, input: LogicalPlan, limit_clause: Option, + fetch_clause: Option, planner_context: &mut PlannerContext, ) -> Result { - let Some(limit_clause) = limit_clause else { + if limit_clause.is_none() && fetch_clause.is_none() { return Ok(input); - }; + } let empty_schema = DFSchema::empty(); - let (skip, fetch, limit_by_exprs) = match limit_clause { - LimitClause::LimitOffset { + let (skip, mut fetch, limit_by_exprs) = match limit_clause { + Some(LimitClause::LimitOffset { limit, offset, limit_by, - } => { + }) => { let skip = offset .map(|o| self.sql_to_expr(o.value, &empty_schema, planner_context)) .transpose()?; @@ -274,15 +277,39 @@ impl SqlToRel<'_, S> { (skip, fetch, limit_by_exprs) } - LimitClause::OffsetCommaLimit { offset, limit } => { + Some(LimitClause::OffsetCommaLimit { offset, limit }) => { let skip = Some(self.sql_to_expr(offset, &empty_schema, planner_context)?); let fetch = Some(self.sql_to_expr(limit, &empty_schema, planner_context)?); (skip, fetch, vec![]) } + None => (None, None, vec![]), }; + if let Some(Fetch { + with_ties, + percent, + quantity, + }) = fetch_clause + { + if with_ties { + return not_impl_err!("FETCH WITH TIES is not supported yet"); + } + if percent { + return not_impl_err!("FETCH PERCENT is not supported yet"); + } + if fetch.is_some() { + return not_impl_err!("LIMIT and FETCH cannot be used together"); + } + fetch = Some(match quantity { + Some(quantity) => { + self.sql_to_expr(quantity, &empty_schema, planner_context)? + } + None => lit(1_i64), + }); + } + if !limit_by_exprs.is_empty() { return not_impl_err!("LIMIT BY clause is not supported yet"); } diff --git a/datafusion/sql/tests/sql_integration.rs b/datafusion/sql/tests/sql_integration.rs index 1fa5ce2d6cfd5..04a903e09d45e 100644 --- a/datafusion/sql/tests/sql_integration.rs +++ b/datafusion/sql/tests/sql_integration.rs @@ -4997,10 +4997,57 @@ fn test_offset_after_limit() { } #[test] -fn fetch_clause_is_not_supported() { - let sql = "SELECT 1 FETCH NEXT 1 ROW ONLY"; +fn fetch_clause() { + let sql = "SELECT id FROM person ORDER BY id OFFSET 3 ROWS FETCH NEXT 5 ROWS ONLY"; + let plan = logical_plan(sql).unwrap(); + assert_snapshot!( + plan, + @r" + Limit: skip=3, fetch=5 + Sort: person.id ASC NULLS LAST + Projection: person.id + TableScan: person + " + ); +} + +#[test] +fn fetch_clause_without_quantity_defaults_to_one() { + let sql = "SELECT id FROM person FETCH FIRST ROW ONLY"; + let plan = logical_plan(sql).unwrap(); + assert_snapshot!( + plan, + @r" + Limit: skip=0, fetch=1 + Projection: person.id + TableScan: person + " + ); +} + +#[test] +fn fetch_clause_applies_to_set_operation() { + let sql = "SELECT 1 AS id UNION ALL SELECT 2 FETCH FIRST 1 ROW ONLY"; + let plan = logical_plan(sql).unwrap(); + assert_snapshot!( + plan, + @r" + Limit: skip=0, fetch=1 + Union + Projection: Int64(1) AS id + EmptyRelation: rows=1 + Projection: Int64(2) + EmptyRelation: rows=1 + " + ); +} + +#[rstest] +#[case("SELECT 1 FETCH FIRST 10 PERCENT ROWS ONLY", "FETCH PERCENT")] +#[case("SELECT 1 ORDER BY 1 FETCH FIRST 1 ROW WITH TIES", "FETCH WITH TIES")] +fn unsupported_fetch_options(#[case] sql: &str, #[case] expected: &str) { let err = logical_plan(sql).unwrap_err(); - assert_contains!(err.to_string(), "FETCH clause is not supported yet"); + assert_contains!(err.to_string(), expected); } #[test] diff --git a/datafusion/sqllogictest/test_files/errors.slt b/datafusion/sqllogictest/test_files/errors.slt index ab934279c32ec..92d550b92384b 100644 --- a/datafusion/sqllogictest/test_files/errors.slt +++ b/datafusion/sqllogictest/test_files/errors.slt @@ -74,9 +74,9 @@ statement error DataFusion error: Error during planning: Unsupported compound id SELECT COUNT(*) FROM way.too.many.namespaces.as.ident.prefixes.aggregate_test_100 -# fetch_clause_not_supported -statement error FETCH clause is not supported yet -SELECT 1 FETCH NEXT 1 ROW ONLY +# fetch_with_ties_not_supported +statement error FETCH WITH TIES is not supported yet +SELECT 1 ORDER BY 1 FETCH NEXT 1 ROW WITH TIES