Skip to content

Commit d55153a

Browse files
authored
Plan struct colon access as get_field (#25373)
## Which issue does this PR close? Closes no issue. ## Rationale for this change `SqlToRel::parse_json_access` represents SQL colon access as a colon binary operator. For an Arrow Struct, that expression currently retains the complete Struct type and reaches an unsupported physical binary operator instead of accessing the requested field. ## What changes are included in this PR? `NestedFunctionPlanner` now lowers a static colon access on a Struct to the existing vectorized `get_field` UDF. Nested paths become nested `get_field` calls, which the existing simplifier can flatten and push toward scan leaves. Non-Struct colon expressions remain available to other planners. The tests cover nested Struct access and the non-Struct fallback. ## Are there any user-facing changes? SQL JSON-style access on typed Struct expressions can now be planned and executed through `get_field`. ## How was this change tested? - `cargo test -p datafusion-functions-nested planner::tests --lib` (2 passed) - `cargo fmt --all -- --check` A package clippy run reached an unrelated pre-existing `from_iter_instead_of_collect` warning in `datafusion/common/src/scalar/mod.rs`; the changed crate passed this same focused clippy command on the DataFusion 55 fork.
1 parent d9a9321 commit d55153a

4 files changed

Lines changed: 221 additions & 31 deletions

File tree

‎datafusion/core/tests/user_defined/expr_planner.rs‎

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
use arrow::array::RecordBatch;
1919
use datafusion::common::test_util::batches_to_string;
2020
use std::sync::Arc;
21+
use std::sync::atomic::{AtomicBool, Ordering};
2122

2223
use datafusion::common::DFSchema;
2324
use datafusion::error::Result;
@@ -33,6 +34,24 @@ use datafusion_expr::planner::{ExprPlanner, PlannerResult, RawBinaryExpr};
3334
#[derive(Debug)]
3435
struct MyCustomPlanner;
3536

37+
#[derive(Debug)]
38+
struct DeferColonPlanner {
39+
colon_seen: Arc<AtomicBool>,
40+
}
41+
42+
impl ExprPlanner for DeferColonPlanner {
43+
fn plan_binary_op(
44+
&self,
45+
expr: RawBinaryExpr,
46+
_schema: &DFSchema,
47+
) -> Result<PlannerResult<RawBinaryExpr>> {
48+
if matches!(&expr.op, BinaryOperator::Custom(op) if op == ":") {
49+
self.colon_seen.store(true, Ordering::Relaxed);
50+
}
51+
Ok(PlannerResult::Original(expr))
52+
}
53+
}
54+
3655
impl ExprPlanner for MyCustomPlanner {
3756
fn plan_binary_op(
3857
&self,
@@ -61,6 +80,13 @@ impl ExprPlanner for MyCustomPlanner {
6180
format!("{} ? {}", expr.left, expr.right),
6281
))))
6382
}
83+
BinaryOperator::Custom(op) if op == ":" => {
84+
Ok(PlannerResult::Planned(Expr::Alias(Alias::new(
85+
Expr::Literal(ScalarValue::Boolean(Some(true)), None),
86+
None::<&str>,
87+
"custom colon",
88+
))))
89+
}
6490
_ => Ok(PlannerResult::Original(expr)),
6591
}
6692
}
@@ -125,3 +151,54 @@ async fn test_question_filter() {
125151
+---+
126152
");
127153
}
154+
155+
#[tokio::test]
156+
async fn test_custom_struct_colon_operator() {
157+
let config =
158+
SessionConfig::new().set_str("datafusion.sql_parser.dialect", "snowflake");
159+
let mut ctx = SessionContext::new_with_config(config);
160+
ctx.register_expr_planner(Arc::new(MyCustomPlanner))
161+
.unwrap();
162+
163+
let actual = ctx
164+
.sql("select {'a': 1}:a;")
165+
.await
166+
.unwrap()
167+
.collect()
168+
.await
169+
.unwrap();
170+
insta::assert_snapshot!(batches_to_string(&actual), @r"
171+
+--------------+
172+
| custom colon |
173+
+--------------+
174+
| true |
175+
+--------------+
176+
");
177+
}
178+
179+
#[tokio::test]
180+
async fn test_deferred_struct_colon_operator() {
181+
let config =
182+
SessionConfig::new().set_str("datafusion.sql_parser.dialect", "snowflake");
183+
let mut ctx = SessionContext::new_with_config(config);
184+
let colon_seen = Arc::new(AtomicBool::new(false));
185+
ctx.register_expr_planner(Arc::new(DeferColonPlanner {
186+
colon_seen: Arc::clone(&colon_seen),
187+
}))
188+
.unwrap();
189+
190+
let actual = ctx
191+
.sql("select {'a': 1}:a;")
192+
.await
193+
.unwrap()
194+
.collect()
195+
.await
196+
.unwrap();
197+
assert!(colon_seen.load(Ordering::Relaxed));
198+
assert_eq!(actual.len(), 1);
199+
assert_eq!(actual[0].num_rows(), 1);
200+
assert_eq!(
201+
ScalarValue::try_from_array(actual[0].column(0).as_ref(), 0).unwrap(),
202+
ScalarValue::Int64(Some(1))
203+
);
204+
}

‎datafusion/sql/src/expr/mod.rs‎

Lines changed: 118 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -24,8 +24,8 @@ use datafusion_expr::planner::{
2424
use sqlparser::ast::{
2525
AccessExpr, BinaryOperator, CastFormat, CastKind, CeilFloorKind,
2626
DataType as SQLDataType, DateTimeField, DictionaryField, Expr as SQLExpr,
27-
ExprWithAlias as SQLExprWithAlias, JsonPath, MapEntry, Spanned, StructField,
28-
Subscript, TrimWhereField, TypedString, Value, ValueWithSpan,
27+
ExprWithAlias as SQLExprWithAlias, JsonPath, JsonPathElem, MapEntry, Spanned,
28+
StructField, Subscript, TrimWhereField, TypedString, Value, ValueWithSpan,
2929
};
3030
use sqlparser::ast::{Query, Visit, Visitor};
3131

@@ -313,21 +313,31 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
313313
right: Expr,
314314
schema: &DFSchema,
315315
) -> Result<Expr> {
316-
// try extension planers
317-
let mut binary_expr = RawBinaryExpr { op, left, right };
316+
let binary_expr = RawBinaryExpr { op, left, right };
317+
match self.try_plan_binary_op(binary_expr, schema)? {
318+
PlannerResult::Planned(expr) => Ok(expr),
319+
PlannerResult::Original(RawBinaryExpr { op, left, right }) => {
320+
self.build_binary_expr(&op, left, right)
321+
}
322+
}
323+
}
324+
325+
fn try_plan_binary_op(
326+
&self,
327+
mut binary_expr: RawBinaryExpr,
328+
schema: &DFSchema,
329+
) -> Result<PlannerResult<RawBinaryExpr>> {
318330
for planner in self.context_provider.get_expr_planners() {
319331
match planner.plan_binary_op(binary_expr, schema)? {
320332
PlannerResult::Planned(expr) => {
321-
return Ok(expr);
333+
return Ok(PlannerResult::Planned(expr));
322334
}
323335
PlannerResult::Original(expr) => {
324336
binary_expr = expr;
325337
}
326338
}
327339
}
328-
329-
let RawBinaryExpr { op, left, right } = binary_expr;
330-
self.build_binary_expr(&op, left, right)
340+
Ok(PlannerResult::Original(binary_expr))
331341
}
332342

333343
pub fn sql_to_expr_with_alias(
@@ -840,20 +850,72 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
840850
value: Box<SQLExpr>,
841851
path: &JsonPath,
842852
) -> Result<Expr> {
853+
let value = self.sql_to_expr(*value, schema, planner_context)?;
843854
let json_path = path.to_string();
844-
let json_path = if let Some(json_path) = json_path.strip_prefix(":") {
845-
// sqlparser's JsonPath display adds an extra `:` at the beginning.
846-
json_path.to_owned()
847-
} else {
848-
json_path
855+
let json_path = json_path.strip_prefix(":").unwrap_or(&json_path);
856+
let binary_expr = RawBinaryExpr {
857+
op: BinaryOperator::Custom(":".to_owned()),
858+
left: value,
859+
right: Expr::Literal(ScalarValue::Utf8(Some(json_path.to_owned())), None),
849860
};
850-
self.build_logical_expr(
851-
BinaryOperator::Custom(":".to_owned()),
852-
self.sql_to_expr(*value, schema, planner_context)?,
853-
// pass json path as a string literal, let the impl parse it when needed.
854-
Expr::Literal(ScalarValue::Utf8(Some(json_path)), None),
855-
schema,
856-
)
861+
let binary_expr = match self.try_plan_binary_op(binary_expr, schema)? {
862+
PlannerResult::Planned(expr) => return Ok(expr),
863+
PlannerResult::Original(expr) => expr,
864+
};
865+
866+
if !path.path.is_empty()
867+
&& matches!(&binary_expr.op, BinaryOperator::Custom(op) if op == ":")
868+
&& is_struct_like(&binary_expr.left.get_type(schema)?)
869+
&& path
870+
.path
871+
.iter()
872+
.all(|element| struct_field_name_from_json_path_elem(element).is_some())
873+
{
874+
let mut planned = binary_expr.left.clone();
875+
let mut all_fields_planned = true;
876+
877+
for element in &path.path {
878+
let field_name = struct_field_name_from_json_path_elem(element)
879+
.expect("all path elements were checked above");
880+
let field_access = RawFieldAccessExpr {
881+
expr: planned,
882+
field_access: GetFieldAccess::NamedStructField {
883+
name: ScalarValue::from(field_name),
884+
},
885+
};
886+
match self.try_plan_field_access(field_access, schema)? {
887+
PlannerResult::Planned(expr) => planned = expr,
888+
PlannerResult::Original(field_access) => {
889+
planned = field_access.expr;
890+
all_fields_planned = false;
891+
break;
892+
}
893+
}
894+
}
895+
896+
if all_fields_planned {
897+
return Ok(planned);
898+
}
899+
}
900+
901+
let RawBinaryExpr { op, left, right } = binary_expr;
902+
self.build_binary_expr(&op, left, right)
903+
}
904+
905+
fn try_plan_field_access(
906+
&self,
907+
mut field_access_expr: RawFieldAccessExpr,
908+
schema: &DFSchema,
909+
) -> Result<PlannerResult<RawFieldAccessExpr>> {
910+
for planner in self.context_provider.get_expr_planners() {
911+
match planner.plan_field_access(field_access_expr, schema)? {
912+
PlannerResult::Planned(expr) => {
913+
return Ok(PlannerResult::Planned(expr));
914+
}
915+
PlannerResult::Original(expr) => field_access_expr = expr,
916+
}
917+
}
918+
Ok(PlannerResult::Original(field_access_expr))
857919
}
858920

859921
/// Parses a struct(..) expression and plans it creation
@@ -1403,22 +1465,47 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
14031465
.into_iter()
14041466
.flatten()
14051467
.try_fold(root, |expr, field_access| {
1406-
let mut field_access_expr = RawFieldAccessExpr { expr, field_access };
1407-
for planner in self.context_provider.get_expr_planners() {
1408-
match planner.plan_field_access(field_access_expr, schema)? {
1409-
PlannerResult::Planned(expr) => return Ok(expr),
1410-
PlannerResult::Original(expr) => {
1411-
field_access_expr = expr;
1412-
}
1413-
}
1468+
let field_access_expr = RawFieldAccessExpr { expr, field_access };
1469+
match self.try_plan_field_access(field_access_expr, schema)? {
1470+
PlannerResult::Planned(expr) => Ok(expr),
1471+
PlannerResult::Original(field_access_expr) => not_impl_err!(
1472+
"GetFieldAccess not supported by ExprPlanner: {field_access_expr:?}"
1473+
),
14141474
}
1415-
not_impl_err!(
1416-
"GetFieldAccess not supported by ExprPlanner: {field_access_expr:?}"
1417-
)
14181475
})
14191476
}
14201477
}
14211478

1479+
fn is_struct_like(data_type: &DataType) -> bool {
1480+
matches!(data_type, DataType::Struct(_))
1481+
|| matches!(
1482+
data_type,
1483+
DataType::Dictionary(_, value_type)
1484+
if matches!(value_type.as_ref(), DataType::Struct(_))
1485+
)
1486+
}
1487+
1488+
fn struct_field_name_from_json_path_elem(element: &JsonPathElem) -> Option<&str> {
1489+
match element {
1490+
JsonPathElem::Dot { key, .. } => Some(key),
1491+
JsonPathElem::Bracket {
1492+
key:
1493+
SQLExpr::Value(ValueWithSpan {
1494+
value: Value::SingleQuotedString(key) | Value::DoubleQuotedString(key),
1495+
span: _,
1496+
}),
1497+
}
1498+
| JsonPathElem::ColonBracket {
1499+
key:
1500+
SQLExpr::Value(ValueWithSpan {
1501+
value: Value::SingleQuotedString(key) | Value::DoubleQuotedString(key),
1502+
span: _,
1503+
}),
1504+
} => Some(key),
1505+
JsonPathElem::Bracket { .. } | JsonPathElem::ColonBracket { .. } => None,
1506+
}
1507+
}
1508+
14221509
/// Builds a CASE expression that handles NULL semantics for `x <op> ANY(arr)`:
14231510
///
14241511
/// ```text

‎datafusion/sqllogictest/test_files/dictionary_struct.slt‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,16 @@ Carol
109109
Alice
110110
Bob
111111

112+
# Snowflake-style colon access preserves dictionary encoding and values.
113+
query TT
114+
SELECT dict_struct:name, arrow_typeof(dict_struct:name) FROM dict_struct_table;
115+
----
116+
Alice Dictionary(UInt32, Utf8)
117+
Bob Dictionary(UInt32, Utf8)
118+
Carol Dictionary(UInt32, Utf8)
119+
Alice Dictionary(UInt32, Utf8)
120+
Bob Dictionary(UInt32, Utf8)
121+
112122
# A nullable parent makes its non-nullable child nullable, preserving the dictionary type.
113123
query TTB
114124
SELECT ds['name'], arrow_typeof(ds['name']), arrow_field(ds['name'])['nullable'] FROM dict_struct_nullable;

‎datafusion/sqllogictest/test_files/struct.slt‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -727,6 +727,22 @@ select get_field(s, 'inner', 'val') from nested_test;
727727
----
728728
100
729729

730+
# Snowflake-style colon access uses the structured JSON path before it is serialized.
731+
query III
732+
select
733+
s:inner.val as unquoted,
734+
s:"inner"."val" as quoted,
735+
s:['inner']['val'] as bracketed
736+
from nested_test;
737+
----
738+
100 100 100
739+
740+
# Quoted dots are part of the field name, not path separators.
741+
query I
742+
select named_struct('inner.val', 200):"inner.val";
743+
----
744+
200
745+
730746
statement ok
731747
drop table nested_test;
732748

0 commit comments

Comments
 (0)