From 8f1bca307d9248c9bec4dfea716112bf61aaf171 Mon Sep 17 00:00:00 2001 From: Jia-Xuan Liu Date: Wed, 27 Nov 2024 10:15:10 +0800 Subject: [PATCH 1/3] upgrade the required sqlparser version --- Cargo.toml | 2 +- datafusion/sql/src/expr/function.rs | 5 ++ datafusion/sql/src/expr/mod.rs | 82 +++++++++++++++++++++++++++-- datafusion/sql/src/statement.rs | 55 +++++++++++-------- datafusion/sql/src/unparser/ast.rs | 1 + datafusion/sql/src/unparser/plan.rs | 9 +++- 6 files changed, 126 insertions(+), 28 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index fa0ce2461fafa..2d5ce1be7a434 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -149,7 +149,7 @@ recursive = "0.1.1" regex = "1.8" rstest = "0.23.0" serde_json = "1" -sqlparser = { version = "0.52.0", features = ["visitor"] } +sqlparser = { git = "https://github.com/goldmedal/sqlparser-rs.git", branch = "wren/0.12.3-array-struct", features = ["visitor"] } tempfile = "3" tokio = { version = "1.36", features = ["macros", "rt", "sync"] } url = "2.2" diff --git a/datafusion/sql/src/expr/function.rs b/datafusion/sql/src/expr/function.rs index cb7255bb7873c..60ab33000740a 100644 --- a/datafusion/sql/src/expr/function.rs +++ b/datafusion/sql/src/expr/function.rs @@ -169,6 +169,11 @@ impl FunctionArgs { "Calling {name}: SEPARATOR not supported in function arguments: {sep}" ) } + FunctionArgumentClause::JsonNullClause(jn) => { + return not_impl_err!( + "Calling {name}: JSON NULL clause not supported in function arguments: {jn}" + ) + } } } diff --git a/datafusion/sql/src/expr/mod.rs b/datafusion/sql/src/expr/mod.rs index a095d895b46fd..edfd62103df75 100644 --- a/datafusion/sql/src/expr/mod.rs +++ b/datafusion/sql/src/expr/mod.rs @@ -23,9 +23,11 @@ use datafusion_expr::planner::{ use recursive::recursive; use sqlparser::ast::{ BinaryOperator, CastFormat, CastKind, DataType as SQLDataType, DictionaryField, - Expr as SQLExpr, MapEntry, StructField, Subscript, TrimWhereField, Value, + Expr as SQLExpr, MapAccessKey, MapAccessSyntax, MapEntry, StructField, Subscript, + TrimWhereField, Value, }; +use crate::planner::{ContextProvider, PlannerContext, SqlToRel}; use datafusion_common::{ internal_datafusion_err, internal_err, not_impl_err, plan_err, DFSchema, Result, ScalarValue, @@ -37,8 +39,6 @@ use datafusion_expr::{ Operator, TryCast, }; -use crate::planner::{ContextProvider, PlannerContext, SqlToRel}; - mod binary_op; mod function; mod grouping_set; @@ -209,8 +209,9 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { self.sql_identifier_to_expr(id, schema, planner_context) } - SQLExpr::MapAccess { .. } => { - not_impl_err!("Map Access") + // . + SQLExpr::MapAccess { column, keys } => { + self.sql_map_access_to_expr(*column, keys, schema, planner_context) } // ["foo"], [4] or [4:5] @@ -1036,6 +1037,77 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { "GetFieldAccess not supported by ExprPlanner: {field_access_expr:?}" ) } + + fn sql_map_access_to_expr( + &self, + expr: SQLExpr, + keys: Vec, + schema: &DFSchema, + planner_context: &mut PlannerContext, + ) -> Result { + let head = + self.sql_expr_to_logical_expr(expr.clone(), schema, planner_context)?; + let data_type = head.get_type(schema)?; + let field_accesses = keys + .iter() + .map(|key| match &key.syntax { + MapAccessSyntax::Bracket => match data_type { + DataType::List(_) + | DataType::FixedSizeList(_, _) + | DataType::LargeList(_) => Ok(GetFieldAccess::ListIndex { + key: Box::new(self.sql_expr_to_logical_expr( + key.key.clone(), + schema, + planner_context, + )?), + }), + DataType::Map(_, _) | DataType::Struct(_) => { + let name = get_string_value(&key.key)?; + Ok(GetFieldAccess::NamedStructField { name }) + } + _ => not_impl_err!( + "MapAccessKey not supported for data type: {data_type}" + ), + }, + MapAccessSyntax::Period => match data_type { + DataType::Map(_, _) | DataType::Struct(_) => { + let name = get_string_value(&key.key)?; + Ok(GetFieldAccess::NamedStructField { name }) + } + _ => not_impl_err!( + "MapAccessKey not supported for data type: {data_type}" + ), + }, + }) + .collect::>>()?; + + field_accesses + .into_iter() + .try_fold(head, |expr, field_access| { + let mut field_access_expr = RawFieldAccessExpr { expr, field_access }; + for planner in self.context_provider.get_expr_planners() { + match planner.plan_field_access(field_access_expr, schema)? { + PlannerResult::Planned(expr) => return Ok(expr), + PlannerResult::Original(expr) => { + field_access_expr = expr; + } + } + } + not_impl_err!( + "MapAccessKey not supported by ExprPlanner: {field_access_expr:?}" + ) + }) + } +} + +fn get_string_value(expr: &SQLExpr) -> Result { + match expr { + SQLExpr::Value(Value::SingleQuotedString(s) | Value::DoubleQuotedString(s)) => { + Ok(ScalarValue::from(s.clone())) + } + SQLExpr::Identifier(id) => Ok(ScalarValue::from(id.value.clone())), + _ => not_impl_err!("Expected string value or identifier, found: {expr:?}"), + } } #[cfg(test)] diff --git a/datafusion/sql/src/statement.rs b/datafusion/sql/src/statement.rs index 31b836f32b242..088c1fd47e14f 100644 --- a/datafusion/sql/src/statement.rs +++ b/datafusion/sql/src/statement.rs @@ -54,13 +54,13 @@ use datafusion_expr::{ TransactionConclusion, TransactionEnd, TransactionIsolationLevel, TransactionStart, Volatility, WriteOp, }; -use sqlparser::ast::{self, SqliteOnConflict}; +use sqlparser::ast::{self, ShowStatementOptions, SqliteOnConflict}; use sqlparser::ast::{ Assignment, AssignmentTarget, ColumnDef, CreateIndex, CreateTable, CreateTableOptions, Delete, DescribeAlias, Expr as SQLExpr, FromTable, Ident, Insert, ObjectName, ObjectType, OneOrManyWithParens, Query, SchemaName, SetExpr, - ShowCreateObject, ShowStatementFilter, Statement, TableConstraint, TableFactor, - TableWithJoins, TransactionMode, UnaryOperator, Value, + ShowCreateObject, Statement, TableConstraint, TableFactor, TableWithJoins, + TransactionMode, UnaryOperator, Value, }; use sqlparser::parser::ParserError::ParserError; @@ -683,21 +683,26 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { ))), Statement::ShowTables { + terse, + history, extended, full, - db_name, - filter, - // SHOW TABLES IN/FROM are equivalent, this field specifies which the user - // specified, but it doesn't affect the plan so ignore the field - clause: _, - } => self.show_tables_to_plan(extended, full, db_name, filter), + external, + show_options, + } => self.show_tables_to_plan( + terse, + history, + extended, + full, + external, + show_options, + ), Statement::ShowColumns { extended, full, - table_name, - filter, - } => self.show_columns_to_plan(extended, full, table_name, filter), + show_options, + } => self.show_columns_to_plan(extended, full, show_options), Statement::Insert(Insert { or, @@ -766,10 +771,14 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { from, selection, returning, + or, } => { if returning.is_some() { plan_err!("Update-returning clause not yet supported")?; } + if or.is_some() { + plan_err!("Update-or clause not yet supported")?; + } self.update_to_plan(table, assignments, from, selection) } @@ -1067,15 +1076,17 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { /// Generate a logical plan from a "SHOW TABLES" query fn show_tables_to_plan( &self, + terse: bool, + history: bool, extended: bool, full: bool, - db_name: Option, - filter: Option, + external: bool, + _show_options: ShowStatementOptions, ) -> Result { if self.has_table("information_schema", "tables") { // We only support the basic "SHOW TABLES" // https://github.com/apache/datafusion/issues/3188 - if db_name.is_some() || filter.is_some() || full || extended { + if terse || history || full || extended || external { plan_err!("Unsupported parameters to SHOW TABLES") } else { let query = "SELECT * FROM information_schema.tables;"; @@ -1841,18 +1852,20 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { &self, extended: bool, full: bool, - sql_table_name: ObjectName, - filter: Option, + show_options: ShowStatementOptions, ) -> Result { - if filter.is_some() { - return plan_err!("SHOW COLUMNS with WHERE or LIKE is not supported"); - } - if !self.has_table("information_schema", "columns") { return plan_err!( "SHOW COLUMNS is not supported unless information_schema is enabled" ); } + if show_options.filter_position.is_some() { + return plan_err!("SHOW COLUMNS with WHERE or LIKE is not supported"); + } + let Some(sql_table_name) = show_options.show_in.and_then(|i| i.parent_name) + else { + return plan_err!("SHOW COLUMNS is not supported without a table name"); + }; // Figure out the where clause let where_clause = object_name_to_qualifier( &sql_table_name, diff --git a/datafusion/sql/src/unparser/ast.rs b/datafusion/sql/src/unparser/ast.rs index cc0812cd71e12..268b4ffd9305e 100644 --- a/datafusion/sql/src/unparser/ast.rs +++ b/datafusion/sql/src/unparser/ast.rs @@ -458,6 +458,7 @@ impl TableRelationBuilder { version: self.version.clone(), partitions: self.partitions.clone(), with_ordinality: false, + json_path: None, }) } fn create_empty() -> Self { diff --git a/datafusion/sql/src/unparser/plan.rs b/datafusion/sql/src/unparser/plan.rs index 81e47ed939f22..367e3ff1ead13 100644 --- a/datafusion/sql/src/unparser/plan.rs +++ b/datafusion/sql/src/unparser/plan.rs @@ -42,7 +42,7 @@ use datafusion_expr::{ expr::Alias, BinaryExpr, Distinct, Expr, JoinConstraint, JoinType, LogicalPlan, LogicalPlanBuilder, Operator, Projection, SortExpr, TableScan, }; -use sqlparser::ast::{self, Ident, SetExpr}; +use sqlparser::ast::{self, Ident, SetExpr, TableAliasColumnDef}; use std::sync::Arc; /// Convert a DataFusion [`LogicalPlan`] to [`ast::Statement`] @@ -1006,6 +1006,13 @@ impl Unparser<'_> { } fn new_table_alias(&self, alias: String, columns: Vec) -> ast::TableAlias { + let columns = columns + .into_iter() + .map(|c| TableAliasColumnDef { + name: c, + data_type: None, + }) + .collect(); ast::TableAlias { name: self.new_ident_quoted_if_needs(alias), columns, From d651064eba78a80c03d5df8cc558936598205781 Mon Sep 17 00:00:00 2001 From: Jia-Xuan Liu Date: Wed, 27 Nov 2024 10:15:30 +0800 Subject: [PATCH 2/3] fix upgrade --- datafusion/sql/src/planner.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/datafusion/sql/src/planner.rs b/datafusion/sql/src/planner.rs index ccb2ccf7126f1..3b67e76d92103 100644 --- a/datafusion/sql/src/planner.rs +++ b/datafusion/sql/src/planner.rs @@ -339,7 +339,12 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> { plan: LogicalPlan, alias: TableAlias, ) -> Result { - let plan = self.apply_expr_alias(plan, alias.columns)?; + let idents = alias + .columns + .iter() + .map(|column| column.name.clone()) + .collect::>(); + let plan = self.apply_expr_alias(plan, idents)?; LogicalPlanBuilder::from(plan) .alias(TableReference::bare( From 1862366b3543b7522e0601a7e29a7bd97adc599c Mon Sep 17 00:00:00 2001 From: Jia-Xuan Liu Date: Wed, 27 Nov 2024 10:31:37 +0800 Subject: [PATCH 3/3] add simple test --- datafusion/sqllogictest/test_files/struct.slt | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/datafusion/sqllogictest/test_files/struct.slt b/datafusion/sqllogictest/test_files/struct.slt index 7596b820c688b..f74d161bba453 100644 --- a/datafusion/sqllogictest/test_files/struct.slt +++ b/datafusion/sqllogictest/test_files/struct.slt @@ -595,3 +595,14 @@ Struct([Field { name: "r", data_type: Utf8, nullable: true, dict_id: 0, dict_is_ statement ok drop table t; + +statement ok +create table t(array_col array) as values (array[row('a', 1), row('b', 2)]); + +query TI +select array_col[1].r, array_col[1].c from t; +---- +a 1 + +statement ok +drop table t;