Skip to content
Closed
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
5 changes: 5 additions & 0 deletions datafusion/sql/src/expr/function.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
)
}
}
}

Expand Down
82 changes: 77 additions & 5 deletions datafusion/sql/src/expr/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -37,8 +39,6 @@ use datafusion_expr::{
Operator, TryCast,
};

use crate::planner::{ContextProvider, PlannerContext, SqlToRel};

mod binary_op;
mod function;
mod grouping_set;
Expand Down Expand Up @@ -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")
// <expr>.<field>
SQLExpr::MapAccess { column, keys } => {
self.sql_map_access_to_expr(*column, keys, schema, planner_context)
}

// <expr>["foo"], <expr>[4] or <expr>[4:5]
Expand Down Expand Up @@ -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<MapAccessKey>,
schema: &DFSchema,
planner_context: &mut PlannerContext,
) -> Result<Expr> {
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::<Result<Vec<_>>>()?;

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<ScalarValue> {
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)]
Expand Down
7 changes: 6 additions & 1 deletion datafusion/sql/src/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -339,7 +339,12 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> {
plan: LogicalPlan,
alias: TableAlias,
) -> Result<LogicalPlan> {
let plan = self.apply_expr_alias(plan, alias.columns)?;
let idents = alias
.columns
.iter()
.map(|column| column.name.clone())
.collect::<Vec<_>>();
let plan = self.apply_expr_alias(plan, idents)?;

LogicalPlanBuilder::from(plan)
.alias(TableReference::bare(
Expand Down
55 changes: 34 additions & 21 deletions datafusion/sql/src/statement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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<Ident>,
filter: Option<ShowStatementFilter>,
external: bool,
_show_options: ShowStatementOptions,
) -> Result<LogicalPlan> {
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;";
Expand Down Expand Up @@ -1841,18 +1852,20 @@ impl<'a, S: ContextProvider> SqlToRel<'a, S> {
&self,
extended: bool,
full: bool,
sql_table_name: ObjectName,
filter: Option<ShowStatementFilter>,
show_options: ShowStatementOptions,
) -> Result<LogicalPlan> {
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,
Expand Down
1 change: 1 addition & 0 deletions datafusion/sql/src/unparser/ast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -458,6 +458,7 @@ impl TableRelationBuilder {
version: self.version.clone(),
partitions: self.partitions.clone(),
with_ordinality: false,
json_path: None,
})
}
fn create_empty() -> Self {
Expand Down
9 changes: 8 additions & 1 deletion datafusion/sql/src/unparser/plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`]
Expand Down Expand Up @@ -1006,6 +1006,13 @@ impl Unparser<'_> {
}

fn new_table_alias(&self, alias: String, columns: Vec<Ident>) -> 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,
Expand Down
11 changes: 11 additions & 0 deletions datafusion/sqllogictest/test_files/struct.slt
Original file line number Diff line number Diff line change
Expand Up @@ -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<struct(r varchar, c int)>) 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;