Skip to content

Commit 907de2e

Browse files
committed
feat!: add reader timezone to map-to-struct
1 parent b69ebcf commit 907de2e

13 files changed

Lines changed: 703 additions & 145 deletions

File tree

datafusion-executor/src/expression.rs

Lines changed: 193 additions & 62 deletions
Large diffs are not rendered by default.

ffi/src/expressions/engine_visitor.rs

Lines changed: 52 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,10 +5,10 @@ use std::ffi::c_void;
55
use delta_kernel::expressions::{
66
ArrayData, BinaryExpression, BinaryExpressionOp, BinaryPredicate, BinaryPredicateOp,
77
ColumnName, Expression, ExpressionRef, ExpressionStructPatch, JunctionPredicate,
8-
JunctionPredicateOp, MapData, MapToStructExpression, OpaqueExpression, OpaqueExpressionOpRef,
9-
OpaquePredicate, OpaquePredicateOpRef, ParseJsonExpression, Predicate, Scalar, StructData,
10-
UnaryExpression, UnaryExpressionOp, UnaryPredicate, UnaryPredicateOp, VariadicExpression,
11-
VariadicExpressionOp,
8+
JunctionPredicateOp, MapData, MapToStructExpression, MapToStructOptions, OpaqueExpression,
9+
OpaqueExpressionOpRef, OpaquePredicate, OpaquePredicateOpRef, ParseJsonExpression, Predicate,
10+
Scalar, StructData, UnaryExpression, UnaryExpressionOp, UnaryPredicate, UnaryPredicateOp,
11+
VariadicExpression, VariadicExpressionOp,
1212
};
1313

1414
use super::kernel_visitor::NullTypeTag;
@@ -174,9 +174,9 @@ pub struct EngineExpressionVisitor {
174174
/// `child_list_id`. The `output_schema` handle specifies the schema to parse the JSON
175175
/// into.
176176
pub visit_parse_json: VisitParseJsonFn,
177-
/// Visits the `MapToStruct` expression belonging to the list identified by `sibling_list_id`.
178-
/// The sub-expression (map column) will be in a _one_ item list identified by `child_list_id`.
179-
/// The output struct schema is determined by the evaluator's result type.
177+
/// Visits a `MapToStruct` expression with default options. The sub-expression is in the
178+
/// one-item list identified by `child_list_id`. Expressions with configured options are
179+
/// reported through `visit_unknown` without visiting the child expression.
180180
pub visit_map_to_struct: VisitUnaryFn,
181181
/// Visits the `LessThan` binary operator belonging to the list identified by
182182
/// `sibling_list_id`. The operands will be in a _two_ item list identified by
@@ -696,7 +696,12 @@ fn visit_expression_impl(
696696
schema_handle
697697
);
698698
}
699-
Expression::MapToStruct(MapToStructExpression { map_expr }) => {
699+
Expression::MapToStruct(map_to_struct)
700+
if map_to_struct.options != MapToStructOptions::default() =>
701+
{
702+
visit_unknown(visitor, sibling_list_id, "map_to_struct")
703+
}
704+
Expression::MapToStruct(MapToStructExpression { map_expr, .. }) => {
700705
let child_list_id = call!(visitor, make_field_list, 1);
701706
visit_expression_impl(visitor, map_expr, child_list_id);
702707
call!(visitor, visit_map_to_struct, sibling_list_id, child_list_id);
@@ -771,7 +776,7 @@ fn visit_predicate_internal(predicate: &Predicate, visitor: &mut EngineExpressio
771776

772777
#[cfg(test)]
773778
mod tests {
774-
use delta_kernel::expressions::{lit, Expression, Scalar};
779+
use delta_kernel::expressions::{lit, Expression, MapToStructOptions, Scalar};
775780
use rstest::rstest;
776781

777782
use super::*;
@@ -791,6 +796,10 @@ mod tests {
791796
sibling_list_id: usize,
792797
parts: Vec<String>,
793798
},
799+
Unknown {
800+
sibling_list_id: usize,
801+
name: String,
802+
},
794803
}
795804

796805
#[derive(Default)]
@@ -847,6 +856,18 @@ mod tests {
847856
});
848857
}
849858

859+
extern "C" fn visit_unknown_name(
860+
data: *mut c_void,
861+
sibling_list_id: usize,
862+
name: KernelStringSlice,
863+
) {
864+
let builder = unsafe { &mut *(data as *mut TestExpressionBuilder) };
865+
let name = unsafe { String::try_from_slice(&name) }.unwrap();
866+
builder.events.push(LiteralEvent::Unknown {
867+
sibling_list_id,
868+
name,
869+
});
870+
}
850871
macro_rules! ignore_fn {
851872
($fn_name:ident $(, $arg_type:ty)*) => {
852873
extern "C" fn $fn_name(
@@ -925,7 +946,7 @@ mod tests {
925946
visit_field_patch: ignore_field_patch,
926947
visit_opaque_expr: ignore_opaque_expr,
927948
visit_opaque_pred: ignore_opaque_pred,
928-
visit_unknown: ignore_string_slice,
949+
visit_unknown: visit_unknown_name,
929950
}
930951
}
931952

@@ -992,4 +1013,25 @@ mod tests {
9921013
assert_eq!(top_level_id, 0);
9931014
assert_eq!(builder.events, vec![expected]);
9941015
}
1016+
1017+
#[test]
1018+
fn timezone_aware_map_to_struct_visits_unknown() {
1019+
let expression = Expression::map_to_struct(
1020+
Expression::column(["partitionValues"]),
1021+
MapToStructOptions::default().with_timestamp_timezone("America/Los_Angeles"),
1022+
);
1023+
let mut builder = TestExpressionBuilder::default();
1024+
let mut visitor = test_visitor(&mut builder);
1025+
1026+
let top_level_id = visit_expression_internal(&expression, &mut visitor);
1027+
1028+
assert_eq!(top_level_id, 0);
1029+
assert_eq!(
1030+
builder.events,
1031+
vec![LiteralEvent::Unknown {
1032+
sibling_list_id: 0,
1033+
name: "map_to_struct".to_string(),
1034+
}]
1035+
);
1036+
}
9951037
}

ffi/src/expressions/kernel_visitor.rs

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ use std::sync::Arc;
66
use delta_kernel::engine::arrow_expression::opaque::ArrowOpaquePredicate;
77
use delta_kernel::expressions::{
88
lit, null_lit, BinaryExpressionOp, BinaryPredicateOp, ColumnName, Expression,
9-
JunctionPredicateOp, Predicate, Scalar, UnaryPredicateOp,
9+
JunctionPredicateOp, MapToStructOptions, Predicate, Scalar, UnaryPredicateOp,
1010
};
1111
use delta_kernel::schema::{DataType, PrimitiveType};
1212
use delta_kernel::DeltaResult;
@@ -692,7 +692,10 @@ pub extern "C" fn visit_expression_map_to_struct(
692692
child_expr: usize,
693693
) -> usize {
694694
unwrap_kernel_expression(state, child_expr).map_or(0, |expr| {
695-
wrap_expression(state, Expression::map_to_struct(expr))
695+
wrap_expression(
696+
state,
697+
Expression::map_to_struct(expr, MapToStructOptions::default()),
698+
)
696699
})
697700
}
698701

ffi/src/test_ffi.rs

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,9 @@ use std::sync::Arc;
66

77
use delta_kernel::expressions::{
88
col, column_name, column_pred, lit, null_lit, ArrayData, BinaryExpressionOp, BinaryPredicateOp,
9-
Expression as Expr, ExpressionStructPatchBuilder, MapData, OpaqueExpressionOp,
10-
OpaquePredicateOp, Predicate as Pred, Scalar, ScalarExpressionEvaluator, StructData,
9+
Expression as Expr, ExpressionStructPatchBuilder, MapData, MapToStructOptions,
10+
OpaqueExpressionOp, OpaquePredicateOp, Predicate as Pred, Scalar, ScalarExpressionEvaluator,
11+
StructData,
1112
};
1213
use delta_kernel::kernel_predicates::{
1314
DirectDataSkippingPredicateEvaluator, DirectPredicateEvaluator,
@@ -158,7 +159,7 @@ pub unsafe extern "C" fn get_testing_kernel_expression() -> Handle<SharedExpress
158159
Expr::struct_from([lit(5_i32), lit(20_i64)]),
159160
Expr::opaque(OpaqueTestOp("foo".to_string()), vec![lit(42), lit(1.111)]),
160161
Expr::unknown("mystery"),
161-
Expr::map_to_struct(col!("pv")),
162+
Expr::map_to_struct(col!("pv"), MapToStructOptions::default()),
162163
Expr::coalesce([col!("col"), lit(0_i32)]),
163164
Expr::array([lit(1_i32), lit(2_i32)]),
164165
];
@@ -254,7 +255,7 @@ pub unsafe extern "C" fn get_simple_testing_kernel_expression() -> Handle<Shared
254255
Expr::binary(BinaryExpressionOp::Multiply, lit(5), lit(6)),
255256
Expr::binary(BinaryExpressionOp::Divide, lit(100), lit(4)),
256257
Expr::struct_from([lit(1_i32), lit(2_i64), lit(3.0_f64)]),
257-
Expr::map_to_struct(col!("partitionValues")),
258+
Expr::map_to_struct(col!("partitionValues"), MapToStructOptions::default()),
258259
];
259260
Arc::new(Expr::struct_from(sub_exprs)).into()
260261
}

kernel/proto/expressions.proto

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,8 +150,16 @@ message ParseJsonExpression {
150150
delta.kernel.schema.StructType output_schema = 2;
151151
}
152152

153+
message MapToStructOptions {
154+
// Optional IANA timezone or fixed offset for offset-less TIMESTAMP values. Explicit offsets win;
155+
// absence means UTC. Ambiguous local times use the earlier instant, nonexistent local times use
156+
// the pre-transition offset, and invalid zones must fail evaluation.
157+
optional string timestamp_timezone = 1;
158+
}
159+
153160
message MapToStructExpression {
154161
Expression map_expr = 1;
162+
MapToStructOptions options = 2;
155163
}
156164

157165
// `nullability_predicate` is optional: when set and it evaluates to false/null, the whole

kernel/src/checkpoint/checkpoint_transform.rs

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
use std::sync::{Arc, LazyLock};
1717

1818
use crate::actions::{ADD_NAME, STATS_PARSED as STATS_PARSED_FIELD};
19-
use crate::expressions::{col, Expression, ExpressionRef, UnaryExpressionOp};
19+
use crate::expressions::{col, Expression, ExpressionRef, MapToStructOptions, UnaryExpressionOp};
2020
use crate::schema::{DataType, SchemaRef, SchemaStructPatchBuilder, StructField, StructType};
2121
use crate::struct_patch::ProjectionStructPatchBuilder;
2222
use crate::table_properties::TableProperties;
@@ -211,7 +211,10 @@ fn build_stats_parsed_expr(stats_schema: &SchemaRef) -> ExpressionRef {
211211
fn build_partition_values_parsed_expr() -> ExpressionRef {
212212
Arc::new(Expression::coalesce([
213213
col!(ADD_NAME, PARTITION_VALUES_PARSED_FIELD),
214-
Expression::map_to_struct(col!(ADD_NAME, PARTITION_VALUES_FIELD)),
214+
Expression::map_to_struct(
215+
col!(ADD_NAME, PARTITION_VALUES_FIELD),
216+
MapToStructOptions::default(),
217+
),
215218
]))
216219
}
217220

0 commit comments

Comments
 (0)