diff --git a/docs/source/contributor-guide/expression-audits/conversion_funcs.md b/docs/source/contributor-guide/expression-audits/conversion_funcs.md index e31d5c92d9b..0b52c96d014 100644 --- a/docs/source/contributor-guide/expression-audits/conversion_funcs.md +++ b/docs/source/contributor-guide/expression-audits/conversion_funcs.md @@ -24,7 +24,7 @@ ## cast - Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8 modulo `Cast.canUpCast` refactored to delegate to `UpCastRule.canUpCast`. -- Spark 3.5.8 (audited 2026-05-27): baseline. `Cast(child, dataType, timeZoneId, evalMode)`; eval modes are `LEGACY`, `ANSI`, `TRY`. The legacy `Cast.canCast` matrix and the `Cast.canAnsiCast` matrix decide acceptance per type pair. Comet routes via `CometCast` (`spark/src/main/scala/org/apache/comet/expressions/CometCast.scala`) using a per-source-type support matrix that returns `Compatible`, `Incompatible(reason)`, or `Unsupported(reason)`; literal children are short-circuited to `Compatible()` so `CometLiteral` validates them. The serialized `Cast` proto carries `datatype`, `evalMode`, `timezone` (default `UTC`), `allowIncompat` (from `spark.comet.expression.Cast.allowIncompatible`), and `isSpark4Plus`. The native side (`native/spark-expr/src/conversion_funcs/cast.rs`) implements explicit per-eval-mode branches for narrowing numeric casts that match Spark's overflow exceptions, and falls through to DataFusion `cast_with_options(safe = !ANSI)` for the rest. +- Spark 3.5.8 (audited 2026-05-27): baseline. `Cast(child, dataType, timeZoneId, evalMode)`; eval modes are `LEGACY`, `ANSI`, `TRY`. The legacy `Cast.canCast` matrix and the `Cast.canAnsiCast` matrix decide acceptance per type pair. Comet routes via `CometCast` (`spark/src/main/scala/org/apache/comet/expressions/CometCast.scala`) using a per-source-type support matrix that returns `Compatible`, `Incompatible(reason)`, or `Unsupported(reason)`; literal children are short-circuited to `Compatible()` so `CometLiteral` validates them. The serialized `Cast` proto carries `datatype`, `evalMode`, `timezone` (default `UTC`), and `isSpark4Plus`; `allowIncompatible` is enforced only by the JVM planner and is not sent to the native cast kernel. The native side (`native/spark-expr/src/conversion_funcs/cast.rs`) implements explicit per-eval-mode branches for narrowing numeric casts that match Spark's overflow exceptions, and falls through to DataFusion `cast_with_options(safe = !ANSI)` for the rest. - Spark 4.0.1 (audited 2026-05-27): `VariantType` added; `StringType` literals replaced with `_: StringType` to accommodate collated strings. `(TimestampType, ByteType|ShortType|IntegerType)` added to `canAnsiCast`. `NullIntolerant` -> `nullIntolerant: Boolean` refactor. New `ToPrettyString.BinaryFormatter` semantics for `Binary -> String` are replicated natively via `spark_binary_formatter`. Numeric-to-numeric matrix unchanged. - Spark 4.1.1 (audited 2026-05-27): `TimeType` added; many `TimeType` arms in `canCast`/`canAnsiCast`. Geospatial `GeographyType` / `GeometryType` types added with their own conversion rules. Numeric-to-numeric matrix unchanged. - Known divergences and gaps: diff --git a/native/core/benches/parquet_timestamp_conversion.rs b/native/core/benches/parquet_timestamp_conversion.rs index d090fc79691..6d6e6209637 100644 --- a/native/core/benches/parquet_timestamp_conversion.rs +++ b/native/core/benches/parquet_timestamp_conversion.rs @@ -119,7 +119,7 @@ fn timestamp_containers() -> [(ArrayRef, DataType); 2] { } fn benchmark(c: &mut Criterion) { - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let mut group = c.benchmark_group("parquet_timestamp_conversion"); for width in [8, 1024] { let (array, target) = array_sibling(width); diff --git a/native/core/src/execution/operators/dynamic_filter/parquet_reader/schema_adapter/tests/resolution.rs b/native/core/src/execution/operators/dynamic_filter/parquet_reader/schema_adapter/tests/resolution.rs index 21b7bcc6315..706f75f7f2e 100644 --- a/native/core/src/execution/operators/dynamic_filter/parquet_reader/schema_adapter/tests/resolution.rs +++ b/native/core/src/execution/operators/dynamic_filter/parquet_reader/schema_adapter/tests/resolution.rs @@ -27,7 +27,7 @@ fn spark_factory( use_field_id: bool, defaults: Option>, ) -> Arc { - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.use_field_id = use_field_id; Arc::new(SparkPhysicalExprAdapterFactory::new(options, defaults)) } diff --git a/native/core/src/execution/operators/iceberg_scan.rs b/native/core/src/execution/operators/iceberg_scan.rs index 9f465329ca0..ad69cac4f03 100644 --- a/native/core/src/execution/operators/iceberg_scan.rs +++ b/native/core/src/execution/operators/iceberg_scan.rs @@ -230,7 +230,7 @@ impl IcebergScanExec { let scan_metrics = scan_result.metrics().clone(); let stream = scan_result.stream(); - let spark_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let spark_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let adapter_factory = SparkPhysicalExprAdapterFactory::new(spark_options, None); let adapted_stream = @@ -630,7 +630,7 @@ mod tests { Field::new("a", DataType::Int64, false), Field::new("b", DataType::Int64, false), ])); - let mut options = super::SparkParquetOptions::new(super::EvalMode::Legacy, "UTC", false); + let mut options = super::SparkParquetOptions::new(super::EvalMode::Legacy, "UTC"); options.case_sensitive = true; let factory = super::SparkPhysicalExprAdapterFactory::new(options, None); for name in ["a", "b"] { diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index b88b477931d..96802ca805d 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -710,7 +710,6 @@ impl PhysicalPlanner { SparkCastOptions::new_with_version( eval_mode, &expr.timezone, - expr.allow_incompat, expr.is_spark4_plus, ), spark_expr.expr_id, @@ -803,7 +802,7 @@ impl PhysicalPlanner { "md5" => Ok(Arc::new(Cast::new( func?, DataType::Utf8, - SparkCastOptions::new_without_timezone(EvalMode::Try, true), + SparkCastOptions::new_without_timezone(EvalMode::Try), None, None, ))), @@ -925,8 +924,7 @@ impl PhysicalPlanner { ))) } ExprStruct::ToPrettyString(expr) => { - let mut spark_cast_options = - SparkCastOptions::new(EvalMode::Try, &expr.timezone, true); + let mut spark_cast_options = SparkCastOptions::new(EvalMode::Try, &expr.timezone); let null_string = "NULL"; spark_cast_options.null_string = null_string.to_string(); spark_cast_options.binary_output_style = @@ -6907,7 +6905,7 @@ mod tests { .with_table_parquet_options(TableParquetOptions::new()), ) as Arc; - let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let expr_adapter_factory: Arc = Arc::new( SparkPhysicalExprAdapterFactory::new(spark_parquet_options, None), diff --git a/native/core/src/parquet/cast_column.rs b/native/core/src/parquet/cast_column.rs index ad7b909768b..3a76e89916f 100644 --- a/native/core/src/parquet/cast_column.rs +++ b/native/core/src/parquet/cast_column.rs @@ -393,7 +393,7 @@ mod tests { let schema = Schema::new(vec![Arc::clone(&input_field)]); let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(struct_arr)]).unwrap(); - let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); opts.use_field_id = true; let col_expr: Arc = Arc::new(Column::new("s", 0)); @@ -466,7 +466,7 @@ mod tests { let expr: Arc = Arc::new(Column::new("ts", 0)); let cast_expr = CometCastColumnExpr::try_new(expr, input_field, target_field, None) .unwrap() - .with_parquet_options(SparkParquetOptions::new(eval_mode, "UTC", false)); + .with_parquet_options(SparkParquetOptions::new(eval_mode, "UTC")); let input = TimestampMillisecondArray::from(vec![Some(1_234), Some(-1_234), None]) .with_timezone_opt(source_tz.clone()); diff --git a/native/core/src/parquet/eager_page_index_reader_factory.rs b/native/core/src/parquet/eager_page_index_reader_factory.rs index dbc7d0f3b80..15bafdd2d63 100644 --- a/native/core/src/parquet/eager_page_index_reader_factory.rs +++ b/native/core/src/parquet/eager_page_index_reader_factory.rs @@ -1341,7 +1341,7 @@ mod tests { false, )])); let mut options = - SparkParquetOptions::new(datafusion_comet_spark_expr::EvalMode::Legacy, "UTC", false); + SparkParquetOptions::new(datafusion_comet_spark_expr::EvalMode::Legacy, "UTC"); options.case_sensitive = false; let factory = EagerPageIndexReaderFactory::new( store, diff --git a/native/core/src/parquet/parquet_exec.rs b/native/core/src/parquet/parquet_exec.rs index 93cffbca67a..fba2d6fc13e 100644 --- a/native/core/src/parquet/parquet_exec.rs +++ b/native/core/src/parquet/parquet_exec.rs @@ -351,8 +351,7 @@ fn get_options( // distinguish INT96-derived TimestampLTZ from a true TimestampNTZ source // and apply the pre-Spark-4 SPARK-36182 rejection (#4219). table_parquet_options.global.coerce_int96_tz = Some("UTC".to_string()); - let mut spark_parquet_options = - SparkParquetOptions::new(EvalMode::Legacy, session_timezone, false); + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, session_timezone); spark_parquet_options.allow_cast_unsigned_ints = true; spark_parquet_options.case_sensitive = case_sensitive; spark_parquet_options.return_null_struct_if_all_fields_missing = diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index e8987a79da6..cc17429a72b 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -87,8 +87,6 @@ pub struct SparkParquetOptions { /// session local timezone by an analyzer in Spark. // TODO we should change timezone to Tz to avoid repeated parsing pub timezone: String, - /// Allow casts that are supported but not guaranteed to be 100% compatible - pub allow_incompat: bool, /// Support casting unsigned ints to signed ints (used by Parquet SchemaAdapter) pub allow_cast_unsigned_ints: bool, /// Whether to read dates/timestamps that were written in the legacy hybrid Julian + Gregorian calendar as it is. If false, throw exceptions instead. If the spark type is TimestampNTZ, this should be true. @@ -121,11 +119,10 @@ pub struct SparkParquetOptions { } impl SparkParquetOptions { - pub fn new(eval_mode: EvalMode, timezone: &str, allow_incompat: bool) -> Self { + pub fn new(eval_mode: EvalMode, timezone: &str) -> Self { Self { eval_mode, timezone: timezone.to_string(), - allow_incompat, allow_cast_unsigned_ints: false, use_legacy_date_timestamp_or_ntz: false, case_sensitive: false, @@ -137,11 +134,10 @@ impl SparkParquetOptions { } } - pub fn new_without_timezone(eval_mode: EvalMode, allow_incompat: bool) -> Self { + pub fn new_without_timezone(eval_mode: EvalMode) -> Self { Self { eval_mode, timezone: "".to_string(), - allow_incompat, allow_cast_unsigned_ints: false, use_legacy_date_timestamp_or_ntz: false, case_sensitive: false, @@ -1405,7 +1401,7 @@ mod tests { let err = spark_parquet_convert( ColumnarValue::Array(Arc::new(array)), &to_type, - &SparkParquetOptions::new(EvalMode::Legacy, "UTC", false), + &SparkParquetOptions::new(EvalMode::Legacy, "UTC"), ) .expect_err("array -> int must be an error"); assert!( @@ -1513,7 +1509,7 @@ mod tests { use datafusion_comet_spark_expr::EvalMode; use std::sync::Arc; - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let overflow_millis = 9_223_372_036_854_776_i64; let millis: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![ Some(overflow_millis), @@ -1633,7 +1629,7 @@ mod tests { )); let target = DataType::Struct(target_fields.into()); let input: ArrayRef = Arc::new(input); - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let output = parquet_convert_array(Arc::clone(&input), &target, &options).unwrap(); assert_eq!(output.data_type(), &target); assert!(output.is_null(0)); @@ -1682,7 +1678,7 @@ mod tests { for timezone in [None::>, Some(Arc::from("UTC"))] { for overflow in [i64::MAX, i64::MIN] { - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let millis: ArrayRef = Arc::new( TimestampMillisecondArray::from(vec![overflow, 7, overflow]) .with_timezone_opt(timezone.clone()), @@ -1822,7 +1818,7 @@ mod tests { use datafusion_comet_spark_expr::EvalMode; use std::sync::Arc; - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); for timezone in [None::>, Some(Arc::from("UTC"))] { for overflow in [i64::MIN, i64::MAX] { let values: ArrayRef = Arc::new( @@ -1918,7 +1914,7 @@ mod tests { use datafusion_comet_spark_expr::EvalMode; use std::sync::Arc; - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let values: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![ i64::MAX, 7, @@ -2125,7 +2121,7 @@ mod tests { )); let to_type = struct_type_with_field_id(vec![("f", DataType::Int32, 1)]); - let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); opts.use_field_id = true; let err = parquet_convert_array(from, &to_type, &opts).unwrap_err(); @@ -2165,7 +2161,7 @@ mod tests { )); let to_type = struct_type_with_field_id(vec![("f", DataType::Int32, 2)]); - let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); opts.use_field_id = true; let result = parquet_convert_array(from, &to_type, &opts).unwrap(); diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index 6bcc8d3518b..495eeb545ec 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -1475,7 +1475,6 @@ impl SparkPhysicalExprAdapter { let mut cast_options = SparkCastOptions::new( self.parquet_options.eval_mode, &self.parquet_options.timezone, - self.parquet_options.allow_incompat, ); cast_options.allow_cast_unsigned_ints = self.parquet_options.allow_cast_unsigned_ints; cast_options.is_adapting_schema = true; @@ -2202,7 +2201,7 @@ pub(crate) mod test { let required_schema = Arc::new(Schema::new(vec![Field::new("col", DataType::Int64, false)])); - let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); spark_parquet_options.allow_type_promotion = false; let expr_adapter_factory: Arc = Arc::new( @@ -2247,7 +2246,7 @@ pub(crate) mod test { let required_schema = Arc::new(Schema::new(vec![Field::new("col", DataType::Int64, false)])); - let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); spark_parquet_options.allow_type_promotion = false; let expr_adapter_factory: Arc = Arc::new( @@ -2375,7 +2374,7 @@ pub(crate) mod test { batch: &RecordBatch, required_schema: SchemaRef, ) -> Result { - let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); spark_parquet_options.allow_cast_unsigned_ints = true; let mut stream = scan_parquet(batch, required_schema, spark_parquet_options)?; stream.next().await.unwrap() @@ -2890,7 +2889,7 @@ pub(crate) mod test { #[test] fn nested_dictionary_containers_are_checked_by_value_type() -> Result<(), DataFusionError> { - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let physical_struct = DataType::Struct(Fields::from(vec![ Field::new("x", DataType::Int64, true), Field::new("unused", DataType::Int32, true), @@ -2925,7 +2924,7 @@ pub(crate) mod test { #[test] fn nested_map_shape_mismatch_is_rejected() -> Result<(), DataFusionError> { - let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let valid = map_type(DataType::Int64); let DataType::Map(entries, _) = &valid else { unreachable!() @@ -3048,7 +3047,7 @@ pub(crate) mod test { let required_schema = struct_schema(vec![ Field::new("b", DataType::Int32, true).with_metadata(id_meta("1")) ]); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.use_field_id = true; let mut stream = scan_parquet(&batch, required_schema, options)?; let err = stream @@ -3075,7 +3074,7 @@ pub(crate) mod test { Arc::new(Int32Array::from(vec![1, 2, 3])), )?; let required_schema = struct_schema(vec![Field::new("x", DataType::Int64, true)]); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.allow_type_promotion = false; let mut stream = scan_parquet(&batch, required_schema, options)?; let err = stream @@ -3101,7 +3100,7 @@ pub(crate) mod test { Arc::new(Int32Array::from(Vec::::new())), )?; let required_schema = struct_schema(vec![Field::new("x", DataType::Int64, true)]); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.allow_type_promotion = false; let mut stream = scan_parquet(&batch, required_schema, options)?; while let Some(batch) = stream.next().await { @@ -3486,7 +3485,7 @@ pub(crate) mod test { Arc::new(Int32Array::from(vec![1, 2, 3])), )?; let required_schema = struct_schema(vec![Field::new("x", DataType::Int64, true)]); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.allow_type_promotion = true; let mut stream = scan_parquet(&batch, required_schema, options)?; let result = stream.next().await.unwrap()?; @@ -3579,7 +3578,7 @@ pub(crate) mod test { Field::new("a", DataType::Int64, nullable), Field::new("b", DataType::Int64, false), ])); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = true; let adapter = SparkPhysicalExprAdapterFactory::new(options, None) .create(logical, physical) @@ -3613,7 +3612,7 @@ pub(crate) mod test { Arc::new(Field::new("dup", DataType::Int64, true)), Arc::new(Field::new("dup", DataType::Int64, true)), ]; - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = true; let error = super::match_struct_fields(&fields, &fields[..1], &options) .expect_err("selected duplicate child must fail") @@ -3648,7 +3647,7 @@ pub(crate) mod test { default.rewrite(Arc::clone(&column)).is_err(), "fixture must reach the default-adapter fallback" ); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = true; let adapter = SparkPhysicalExprAdapterFactory::new(options, None) .create(Arc::clone(&logical), Arc::clone(&physical)) @@ -3714,7 +3713,7 @@ pub(crate) mod test { if mode == "missing" { requested.push(Field::new("missing", DataType::Int64, true)); } - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = true; let mut stream = scan_parquet(&batch, struct_schema(requested), options).unwrap(); let error = stream @@ -3816,7 +3815,7 @@ pub(crate) mod test { reader.schema().field(0).data_type(), required.field(0).data_type() ); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = true; let mut stream = scan_parquet(&batch, required, options).unwrap(); let error = stream @@ -3859,7 +3858,7 @@ pub(crate) mod test { // Read with case-insensitive mode, requesting column "b" which matches both "B" and "b" let required_schema = Arc::new(Schema::new(vec![Field::new("b", DataType::Int32, false)])); - let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); spark_parquet_options.case_sensitive = false; let expr_adapter_factory: Arc = Arc::new( @@ -3919,7 +3918,7 @@ pub(crate) mod test { None, )?)); let defaults = HashMap::from([(Column::new("missing", 0), default)]); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = false; let adapter = SparkPhysicalExprAdapterFactory::new(options, Some(defaults)) .create(logical, physical)?; @@ -4014,7 +4013,7 @@ pub(crate) mod test { Field::new("ω", DataType::Int32, true).with_metadata(id_meta("2")), ])); - let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); opts.case_sensitive = false; opts.use_field_id = true; let adapter = SparkPhysicalExprAdapterFactory::new(opts, None) @@ -4045,7 +4044,7 @@ pub(crate) mod test { let physical = Arc::new(Schema::new(vec![ Field::new("MÜNCHEN", storage, true).with_extension_type(VariantType) ])); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = false; let adapter = SparkPhysicalExprAdapterFactory::new(options, None) .create(Arc::clone(&logical), Arc::clone(&physical)) @@ -4104,7 +4103,7 @@ pub(crate) mod test { .is_some()); } let adapter = SparkPhysicalExprAdapterFactory::new( - SparkParquetOptions::new(EvalMode::Legacy, "UTC", false), + SparkParquetOptions::new(EvalMode::Legacy, "UTC"), None, ) .create(logical, physical) @@ -4140,7 +4139,7 @@ pub(crate) mod test { Field::new("other", storage.clone(), true).with_metadata(id_meta("1")), Field::new("__comet_unmatched_field_id_1", DataType::Binary, true), ])); - let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); options.case_sensitive = case_sensitive; options.use_field_id = true; let adapter = SparkPhysicalExprAdapterFactory::new(options, None) @@ -4210,7 +4209,7 @@ pub(crate) mod test { )); let target_type = DataType::List(Arc::clone(&to_item_field)); - let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); let comet_result = spark_parquet_convert( ColumnarValue::Array(Arc::clone(&list_array)), &target_type, @@ -4321,7 +4320,7 @@ pub(crate) mod test { } fn default_options() -> SparkParquetOptions { - SparkParquetOptions::new(EvalMode::Legacy, "UTC", false) + SparkParquetOptions::new(EvalMode::Legacy, "UTC") } /// Dropping a struct field by exact name, including through nested struct-in-struct and @@ -4683,7 +4682,7 @@ pub(crate) mod test { )) as Arc; let batch = RecordBatch::try_new(file_schema, vec![col]).unwrap(); - let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC"); opts.use_field_id = true; let err = match scan_parquet(&batch, required_schema, opts) { diff --git a/native/proto/src/proto/expr.proto b/native/proto/src/proto/expr.proto index d1bdd9c80a5..b85f08f6531 100644 --- a/native/proto/src/proto/expr.proto +++ b/native/proto/src/proto/expr.proto @@ -434,7 +434,8 @@ message Cast { DataType datatype = 2; string timezone = 3; EvalMode eval_mode = 4; - bool allow_incompat = 5; + reserved 5; + reserved "allow_incompat"; // True when running against Spark 4.0+. Controls version-specific cast behaviour // such as the handling of leading whitespace before T-prefixed time-only strings. bool is_spark4_plus = 6; diff --git a/native/spark-expr/benches/cast_binary_to_string.rs b/native/spark-expr/benches/cast_binary_to_string.rs index 8ae9b617615..7c3e5d3a9c5 100644 --- a/native/spark-expr/benches/cast_binary_to_string.rs +++ b/native/spark-expr/benches/cast_binary_to_string.rs @@ -46,7 +46,7 @@ fn build(size: usize, width: usize, null_every: usize) -> ArrayRef { } fn options(style: Option) -> SparkCastOptions { - let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); options.binary_output_style = style; options } diff --git a/native/spark-expr/benches/cast_decimal_to_boolean.rs b/native/spark-expr/benches/cast_decimal_to_boolean.rs index 0d1bfe255c3..2e5af4809e2 100644 --- a/native/spark-expr/benches/cast_decimal_to_boolean.rs +++ b/native/spark-expr/benches/cast_decimal_to_boolean.rs @@ -59,7 +59,7 @@ fn criterion_benchmark(c: &mut Criterion) { let cast_to_bool = Cast::new( expr, DataType::Boolean, - SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + SparkCastOptions::new(EvalMode::Legacy, "UTC"), None, None, ); diff --git a/native/spark-expr/benches/cast_decimal_to_string.rs b/native/spark-expr/benches/cast_decimal_to_string.rs index d0629300a97..5fd76e44fd4 100644 --- a/native/spark-expr/benches/cast_decimal_to_string.rs +++ b/native/spark-expr/benches/cast_decimal_to_string.rs @@ -43,7 +43,7 @@ fn cast_to_utf8() -> Cast { Cast::new( Arc::new(Column::new("a", 0)), DataType::Utf8, - SparkCastOptions::new_without_timezone(EvalMode::Legacy, false), + SparkCastOptions::new_without_timezone(EvalMode::Legacy), None, None, ) diff --git a/native/spark-expr/benches/cast_float_to_decimal.rs b/native/spark-expr/benches/cast_float_to_decimal.rs index fab3c3d1122..1bb071092d3 100644 --- a/native/spark-expr/benches/cast_float_to_decimal.rs +++ b/native/spark-expr/benches/cast_float_to_decimal.rs @@ -52,7 +52,7 @@ fn cast(to: DataType, mode: EvalMode) -> Cast { Cast::new( Arc::new(Column::new("a", 0)), to, - SparkCastOptions::new_without_timezone(mode, false), + SparkCastOptions::new_without_timezone(mode), None, None, ) diff --git a/native/spark-expr/benches/cast_float_to_string.rs b/native/spark-expr/benches/cast_float_to_string.rs index deb24a3e04c..71af3074c8b 100644 --- a/native/spark-expr/benches/cast_float_to_string.rs +++ b/native/spark-expr/benches/cast_float_to_string.rs @@ -67,7 +67,7 @@ fn criterion_benchmark(c: &mut Criterion) { let size = 8192; let f64_array = create_f64_array(size); let f32_array = create_f32_array(size); - let cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy, false); + let cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy); let mut group = c.benchmark_group("cast_float_to_string"); group.bench_function("cast_f64_to_utf8", |b| { diff --git a/native/spark-expr/benches/cast_from_boolean.rs b/native/spark-expr/benches/cast_from_boolean.rs index caccd67e26d..48f080568aa 100644 --- a/native/spark-expr/benches/cast_from_boolean.rs +++ b/native/spark-expr/benches/cast_from_boolean.rs @@ -26,7 +26,7 @@ use std::sync::Arc; fn criterion_benchmark(c: &mut Criterion) { let expr = Arc::new(Column::new("a", 0)); let boolean_batch = create_boolean_batch(); - let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let cast_to_i8 = Cast::new( expr.clone(), DataType::Int8, diff --git a/native/spark-expr/benches/cast_from_string.rs b/native/spark-expr/benches/cast_from_string.rs index 31cf11a84ba..cd344670127 100644 --- a/native/spark-expr/benches/cast_from_string.rs +++ b/native/spark-expr/benches/cast_from_string.rs @@ -41,7 +41,7 @@ fn criterion_benchmark(c: &mut Criterion) { let expr = Arc::new(Column::new("a", 0)); for (mode, mode_name) in EVAL_MODES { - let spark_cast_options = SparkCastOptions::new(mode, "", false); + let spark_cast_options = SparkCastOptions::new(mode, ""); let cast_to_i8 = Cast::new( expr.clone(), DataType::Int8, @@ -88,7 +88,7 @@ fn criterion_benchmark(c: &mut Criterion) { } // Benchmark decimal truncation (Legacy mode only) - let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "", false); + let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, ""); let cast_to_i32 = Cast::new( expr.clone(), DataType::Int32, @@ -116,7 +116,7 @@ fn criterion_benchmark(c: &mut Criterion) { // str -> decimal benchmark let decimal_string_batch = create_decimal_cast_string_batch(); for (mode, mode_name) in EVAL_MODES { - let spark_cast_options = SparkCastOptions::new(mode, "", false); + let spark_cast_options = SparkCastOptions::new(mode, ""); let mut group = c.benchmark_group(format!("cast_string_to_decimal/{}", mode_name)); for (data_type, name) in [ (DataType::Decimal128(38, 10), "decimal_38_10"), @@ -143,7 +143,7 @@ fn criterion_benchmark(c: &mut Criterion) { let float_batch = create_float_string_batch(false); let float_padded_batch = create_float_string_batch(true); for (mode, mode_name) in EVAL_MODES { - let spark_cast_options = SparkCastOptions::new(mode, "", false); + let spark_cast_options = SparkCastOptions::new(mode, ""); let mut group = c.benchmark_group(format!("cast_string_to_bool_and_float/{}", mode_name)); for (data_type, name, batch) in [ (DataType::Boolean, "boolean", &bool_batch), diff --git a/native/spark-expr/benches/cast_int_to_decimal.rs b/native/spark-expr/benches/cast_int_to_decimal.rs index 8949c5c3276..a27b5d277b0 100644 --- a/native/spark-expr/benches/cast_int_to_decimal.rs +++ b/native/spark-expr/benches/cast_int_to_decimal.rs @@ -56,7 +56,7 @@ fn cast(col: &str, to: DataType, mode: EvalMode) -> Cast { Cast::new( Arc::new(Column::new(col, 0)), to, - SparkCastOptions::new_without_timezone(mode, false), + SparkCastOptions::new_without_timezone(mode), None, None, ) diff --git a/native/spark-expr/benches/cast_int_to_timestamp.rs b/native/spark-expr/benches/cast_int_to_timestamp.rs index 4479627ae66..04a62280d5a 100644 --- a/native/spark-expr/benches/cast_int_to_timestamp.rs +++ b/native/spark-expr/benches/cast_int_to_timestamp.rs @@ -27,7 +27,7 @@ const BATCH_SIZE: usize = 8192; fn criterion_benchmark(c: &mut Criterion) { // Test with UTC timezone - let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let timestamp_type = DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())); let mut group = c.benchmark_group("cast_int_to_timestamp"); diff --git a/native/spark-expr/benches/cast_nested.rs b/native/spark-expr/benches/cast_nested.rs index 5c6efd5cb47..0474f2b9af2 100644 --- a/native/spark-expr/benches/cast_nested.rs +++ b/native/spark-expr/benches/cast_nested.rs @@ -93,7 +93,7 @@ fn list_input() -> ArrayRef { } fn criterion_benchmark(c: &mut Criterion) { - let options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let col: Arc = Arc::new(Column::new("c", 0)); let mut group = c.benchmark_group("cast_nested"); diff --git a/native/spark-expr/benches/cast_non_int_numeric_timestamp.rs b/native/spark-expr/benches/cast_non_int_numeric_timestamp.rs index 18ed3041eb6..9c071aae4b6 100644 --- a/native/spark-expr/benches/cast_non_int_numeric_timestamp.rs +++ b/native/spark-expr/benches/cast_non_int_numeric_timestamp.rs @@ -26,7 +26,7 @@ use std::sync::Arc; const BATCH_SIZE: usize = 8192; fn criterion_benchmark(c: &mut Criterion) { - let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let timestamp_type = DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())); let mut group = c.benchmark_group("cast_non_int_numeric_to_timestamp"); diff --git a/native/spark-expr/benches/cast_numeric.rs b/native/spark-expr/benches/cast_numeric.rs index 5153fb7d011..41fd6006d83 100644 --- a/native/spark-expr/benches/cast_numeric.rs +++ b/native/spark-expr/benches/cast_numeric.rs @@ -28,7 +28,7 @@ const NUM_ROWS: usize = 8192; fn criterion_benchmark(c: &mut Criterion) { let batch = create_int32_batch(); let expr = Arc::new(Column::new("a", 0)); - let spark_cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy, false); + let spark_cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy); let cast_i32_to_i8 = Cast::new( expr.clone(), DataType::Int8, @@ -61,7 +61,7 @@ fn criterion_benchmark(c: &mut Criterion) { Cast::new( Arc::new(Column::new("a", 0)), data_type, - SparkCastOptions::new_without_timezone(EvalMode::Legacy, false), + SparkCastOptions::new_without_timezone(EvalMode::Legacy), None, None, ) diff --git a/native/spark-expr/benches/cast_string_to_date.rs b/native/spark-expr/benches/cast_string_to_date.rs index aee0fc0b03a..b2a9d7dabc5 100644 --- a/native/spark-expr/benches/cast_string_to_date.rs +++ b/native/spark-expr/benches/cast_string_to_date.rs @@ -32,7 +32,7 @@ fn criterion_benchmark(c: &mut Criterion) { let cast_to_date = Cast::new( expr, DataType::Date32, - SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + SparkCastOptions::new(EvalMode::Legacy, "UTC"), None, None, ); diff --git a/native/spark-expr/benches/cast_string_to_timestamp.rs b/native/spark-expr/benches/cast_string_to_timestamp.rs index d499cec96ae..117c00592e5 100644 --- a/native/spark-expr/benches/cast_string_to_timestamp.rs +++ b/native/spark-expr/benches/cast_string_to_timestamp.rs @@ -218,7 +218,7 @@ fn criterion_benchmark(c: &mut Criterion) { let cast = Cast::new( Arc::clone(&expr), to_type.clone(), - SparkCastOptions::new(mode, timezone, false), + SparkCastOptions::new(mode, timezone), None, None, ); @@ -238,7 +238,7 @@ fn criterion_benchmark(c: &mut Criterion) { let cast = Cast::new( Arc::clone(&expr), DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())), - SparkCastOptions::new_with_version(EvalMode::Legacy, "UTC", false, true), + SparkCastOptions::new_with_version(EvalMode::Legacy, "UTC", true), None, None, ); diff --git a/native/spark-expr/benches/conditional.rs b/native/spark-expr/benches/conditional.rs index 8af04c2d52f..2febcddee99 100644 --- a/native/spark-expr/benches/conditional.rs +++ b/native/spark-expr/benches/conditional.rs @@ -122,7 +122,7 @@ fn spark_cast(child: Expr, data_type: DataType) -> Expr { Arc::new(Cast::new( child, data_type, - SparkCastOptions::new(EvalMode::Legacy, "America/Los_Angeles", false), + SparkCastOptions::new(EvalMode::Legacy, "America/Los_Angeles"), None, None, )) diff --git a/native/spark-expr/benches/to_csv.rs b/native/spark-expr/benches/to_csv.rs index 8620dd0f160..f51dde2a5d0 100644 --- a/native/spark-expr/benches/to_csv.rs +++ b/native/spark-expr/benches/to_csv.rs @@ -86,7 +86,7 @@ fn criterion_benchmark(c: &mut Criterion) { let default_null_value = ""; let default_quote = "\""; let default_escape = "\\"; - let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, timezone, false); + let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, timezone); cast_options.null_string = default_null_value.to_string(); let csv_write_options = CsvWriteOptions::new( default_delimiter.to_string(), diff --git a/native/spark-expr/benches/wide_decimal.rs b/native/spark-expr/benches/wide_decimal.rs index b38ddbfaa06..6de9fb5e83b 100644 --- a/native/spark-expr/benches/wide_decimal.rs +++ b/native/spark-expr/benches/wide_decimal.rs @@ -85,7 +85,7 @@ fn build_old_expr( ) -> Arc { let left_col: Arc = Arc::new(Column::new("left", 0)); let right_col: Arc = Arc::new(Column::new("right", 1)); - let cast_opts = SparkCastOptions::new_without_timezone(EvalMode::Legacy, false); + let cast_opts = SparkCastOptions::new_without_timezone(EvalMode::Legacy); let left_cast = Arc::new(Cast::new( left_col, DataType::Decimal256(p1, s1), diff --git a/native/spark-expr/src/conditional_funcs/case_when.rs b/native/spark-expr/src/conditional_funcs/case_when.rs index 7dec89c61c0..74df26c977c 100644 --- a/native/spark-expr/src/conditional_funcs/case_when.rs +++ b/native/spark-expr/src/conditional_funcs/case_when.rs @@ -152,7 +152,7 @@ fn coerce_branch( if data_type == common_type { return expr; } - let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); Arc::new(Cast::new( expr, common_type.clone(), @@ -1164,7 +1164,7 @@ mod tests { Arc::new(Cast::new( e, to, - SparkCastOptions::new_without_timezone(EvalMode::Ansi, false), + SparkCastOptions::new_without_timezone(EvalMode::Ansi), None, None, )) diff --git a/native/spark-expr/src/conversion_funcs/boolean.rs b/native/spark-expr/src/conversion_funcs/boolean.rs index 1db2746ce24..7e9dffd3307 100644 --- a/native/spark-expr/src/conversion_funcs/boolean.rs +++ b/native/spark-expr/src/conversion_funcs/boolean.rs @@ -67,7 +67,7 @@ mod tests { } fn test_input_spark_opts() -> SparkCastOptions { - SparkCastOptions::new(EvalMode::Legacy, "Asia/Kolkata", false) + SparkCastOptions::new(EvalMode::Legacy, "Asia/Kolkata") } #[test] diff --git a/native/spark-expr/src/conversion_funcs/cast.rs b/native/spark-expr/src/conversion_funcs/cast.rs index 26659f7489f..30b3fe997a7 100644 --- a/native/spark-expr/src/conversion_funcs/cast.rs +++ b/native/spark-expr/src/conversion_funcs/cast.rs @@ -149,8 +149,6 @@ pub struct SparkCastOptions { /// session local timezone by an analyzer in Spark. // TODO we should change timezone to Tz to avoid repeated parsing pub timezone: String, - /// Allow casts that are supported but not guaranteed to be 100% compatible - pub allow_incompat: bool, /// True when running against Spark 4.0+. Enables version-specific cast behaviour /// such as the handling of leading whitespace before T-prefixed time-only strings. pub is_spark4_plus: bool, @@ -166,11 +164,10 @@ pub struct SparkCastOptions { } impl SparkCastOptions { - pub fn new(eval_mode: EvalMode, timezone: &str, allow_incompat: bool) -> Self { + pub fn new(eval_mode: EvalMode, timezone: &str) -> Self { Self { eval_mode, timezone: timezone.to_string(), - allow_incompat, is_spark4_plus: false, allow_cast_unsigned_ints: false, is_adapting_schema: false, @@ -179,11 +176,10 @@ impl SparkCastOptions { } } - pub fn new_without_timezone(eval_mode: EvalMode, allow_incompat: bool) -> Self { + pub fn new_without_timezone(eval_mode: EvalMode) -> Self { Self { eval_mode, timezone: "".to_string(), - allow_incompat, is_spark4_plus: false, allow_cast_unsigned_ints: false, is_adapting_schema: false, @@ -192,15 +188,10 @@ impl SparkCastOptions { } } - pub fn new_with_version( - eval_mode: EvalMode, - timezone: &str, - allow_incompat: bool, - is_spark4_plus: bool, - ) -> Self { + pub fn new_with_version(eval_mode: EvalMode, timezone: &str, is_spark4_plus: bool) -> Self { Self { is_spark4_plus, - ..Self::new(eval_mode, timezone, allow_incompat) + ..Self::new(eval_mode, timezone) } } } @@ -941,7 +932,7 @@ mod tests { let error = cast_array( Arc::new(StringArray::from(vec!["a"])), &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap_err(); @@ -979,7 +970,7 @@ mod tests { let error = cast_array( input, &output_type, - &SparkCastOptions::new_without_timezone(EvalMode::Ansi, false), + &SparkCastOptions::new_without_timezone(EvalMode::Ansi), ) .unwrap_err(); @@ -1014,7 +1005,7 @@ mod tests { Some(&[0xEDu8, 0xA0, 0x80][..]), ]); // binary_output_style defaults to None, i.e. the plain (non-ToPrettyString) cast path. - let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let result = cast_binary_to_string::(&input, &cast_options).unwrap(); @@ -1038,7 +1029,7 @@ mod tests { Some("abc".as_bytes()), Some(&[0xEDu8, 0xA0, 0x80][..]), ]); - let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); cast_options.binary_output_style = Some(BinaryOutputStyle::Utf8); let result = cast_binary_to_string::(&input, &cast_options).unwrap(); @@ -1060,7 +1051,7 @@ mod tests { Some(b"hi".as_slice()), ])); let cast = |style: Option| { - let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); options.binary_output_style = style; let result = spark_cast( ColumnarValue::Array(Arc::clone(&input)), @@ -1124,7 +1115,7 @@ mod tests { None, Some("héllo".as_bytes()), ])); - let options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let result = spark_cast( ColumnarValue::Array(Arc::clone(&input)), &DataType::Utf8, @@ -1143,7 +1134,7 @@ mod tests { fn test_cast_unsupported_timestamp_to_date() { // Since datafusion uses chrono::Datetime internally not all dates representable by TimestampMicrosecondType are supported let timestamps: PrimitiveArray = vec![i64::MAX].into(); - let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let result = cast_array( Arc::new(timestamps.with_timezone("Europe/Copenhagen")), &DataType::Date32, @@ -1155,7 +1146,7 @@ mod tests { #[test] fn test_cast_invalid_timezone() { let timestamps: PrimitiveArray = vec![i64::MAX].into(); - let cast_options = SparkCastOptions::new(EvalMode::Legacy, "Not a valid timezone", false); + let cast_options = SparkCastOptions::new(EvalMode::Legacy, "Not a valid timezone"); let result = cast_array( Arc::new(timestamps.with_timezone("Europe/Copenhagen")), &DataType::Date32, @@ -1181,7 +1172,7 @@ mod tests { let string_array = cast_array( c, &DataType::Utf8, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap(); let string_array = string_array.as_string::(); @@ -1215,7 +1206,7 @@ mod tests { let cast_array = spark_cast( ColumnarValue::Array(c), &DataType::Struct(fields), - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap(); if let ColumnarValue::Array(cast_array) = cast_array { @@ -1253,7 +1244,7 @@ mod tests { let result = spark_cast( ColumnarValue::Array(outer), &DataType::Struct(to_fields), - &SparkCastOptions::new(EvalMode::Ansi, "UTC", false), + &SparkCastOptions::new(EvalMode::Ansi, "UTC"), ); assert!(result.is_err()); @@ -1278,7 +1269,7 @@ mod tests { let cast_array = spark_cast( ColumnarValue::Array(c), &DataType::Struct(fields), - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap(); if let ColumnarValue::Array(cast_array) = cast_array { @@ -1304,11 +1295,9 @@ mod tests { Arc::new(values_array), None, )); - let string_array = cast_array_to_string( - &list_array, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), - ) - .unwrap(); + let string_array = + cast_array_to_string(&list_array, &SparkCastOptions::new(EvalMode::Legacy, "UTC")) + .unwrap(); let string_array = string_array.as_string::(); assert_eq!(r#"[a, b, c]"#, string_array.value(0)); assert_eq!(r#"[a, null]"#, string_array.value(1)); @@ -1327,11 +1316,9 @@ mod tests { Arc::new(values_array), None, )); - let string_array = cast_array_to_string( - &list_array, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), - ) - .unwrap(); + let string_array = + cast_array_to_string(&list_array, &SparkCastOptions::new(EvalMode::Legacy, "UTC")) + .unwrap(); let string_array = string_array.as_string::(); assert_eq!(r#"[1, 2, 3]"#, string_array.value(0)); assert_eq!(r#"[1, null]"#, string_array.value(1)); @@ -1354,7 +1341,7 @@ mod tests { let to_array = cast_array( from_array, &to_type, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap(); @@ -1368,7 +1355,7 @@ mod tests { assert!(values.iter().all(|value| value.is_none())); } fn legacy_opts() -> SparkCastOptions { - SparkCastOptions::new(EvalMode::Legacy, "UTC", false) + SparkCastOptions::new(EvalMode::Legacy, "UTC") } /// Build a `Map` MapArray (Parquet-style "key_value" field names). @@ -1532,7 +1519,7 @@ mod tests { let casted = cast_array( map_array, &to_type, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), ) .unwrap(); diff --git a/native/spark-expr/src/conversion_funcs/numeric.rs b/native/spark-expr/src/conversion_funcs/numeric.rs index f7fbae54f98..aad9006937c 100644 --- a/native/spark-expr/src/conversion_funcs/numeric.rs +++ b/native/spark-expr/src/conversion_funcs/numeric.rs @@ -2291,7 +2291,7 @@ mod tests { ); for eval_mode in [EvalMode::Legacy, EvalMode::Ansi, EvalMode::Try] { - let options = SparkCastOptions::new(eval_mode, "UTC", false); + let options = SparkCastOptions::new(eval_mode, "UTC"); let doubles = cast_array(Arc::clone(&decimals_38_18), &DataType::Float64, &options).unwrap(); diff --git a/native/spark-expr/src/conversion_funcs/string.rs b/native/spark-expr/src/conversion_funcs/string.rs index adedd3ba0d2..f278bafda68 100644 --- a/native/spark-expr/src/conversion_funcs/string.rs +++ b/native/spark-expr/src/conversion_funcs/string.rs @@ -2491,7 +2491,7 @@ mod tests { for (position, input, expect_value) in cases { for eval_mode in [EvalMode::Legacy, EvalMode::Try, EvalMode::Ansi] { let array: ArrayRef = Arc::new(StringArray::from(vec![Some(input.as_str())])); - let options = SparkCastOptions::new(eval_mode, "UTC", false); + let options = SparkCastOptions::new(eval_mode, "UTC"); let result = cast_array(array, to_type, &options); let context = format!("cast {input:?} ({position}) to {to_type} in {eval_mode:?}"); if expect_value { @@ -2736,7 +2736,7 @@ mod tests { let timezone = "UTC".to_string(); // test casting string dictionary array to timestamp array - let cast_options = SparkCastOptions::new(EvalMode::Legacy, &timezone, false); + let cast_options = SparkCastOptions::new(EvalMode::Legacy, &timezone); let result = cast_array( dict_array, &DataType::Timestamp(TimeUnit::Microsecond, Some(timezone.clone().into())), diff --git a/native/spark-expr/src/conversion_funcs/temporal.rs b/native/spark-expr/src/conversion_funcs/temporal.rs index 31e63644fd5..15926047f73 100644 --- a/native/spark-expr/src/conversion_funcs/temporal.rs +++ b/native/spark-expr/src/conversion_funcs/temporal.rs @@ -121,7 +121,7 @@ mod tests { let target_tz: Option> = Some("UTC".into()); let result = cast_date_to_timestamp( &dates, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), &target_tz, ) .unwrap(); @@ -134,7 +134,7 @@ mod tests { // validate LA timezone (follows Daylight savings) let result = cast_date_to_timestamp( &dates, - &SparkCastOptions::new(EvalMode::Legacy, "America/Los_Angeles", false), + &SparkCastOptions::new(EvalMode::Legacy, "America/Los_Angeles"), &target_tz, ) .unwrap(); @@ -148,7 +148,7 @@ mod tests { // Phoenix timezone (does not follow Daylight savings) let result = cast_date_to_timestamp( &dates, - &SparkCastOptions::new(EvalMode::Legacy, "America/Phoenix", false), + &SparkCastOptions::new(EvalMode::Legacy, "America/Phoenix"), &target_tz, ) .unwrap(); @@ -189,7 +189,7 @@ mod tests { ] { let result = cast_date_to_timestamp( &dates, - &SparkCastOptions::new(EvalMode::Legacy, tz, false), + &SparkCastOptions::new(EvalMode::Legacy, tz), &ntz_target, ) .unwrap(); @@ -215,7 +215,7 @@ mod tests { assert!( cast_date_to_timestamp( &dates, - &SparkCastOptions::new(EvalMode::Legacy, "UTC", false), + &SparkCastOptions::new(EvalMode::Legacy, "UTC"), &ntz_target, ) .is_err(), @@ -240,7 +240,7 @@ mod tests { Some(-106_751_992), ]; for mode in [EvalMode::Legacy, EvalMode::Ansi, EvalMode::Try] { - let options = SparkCastOptions::new(mode, "America/Los_Angeles", false); + let options = SparkCastOptions::new(mode, "America/Los_Angeles"); let result = spark_cast( ColumnarValue::Array(Arc::new(Date32Array::from(days.clone()))), &target, diff --git a/native/spark-expr/src/csv_funcs/to_csv.rs b/native/spark-expr/src/csv_funcs/to_csv.rs index 01fdc901cb7..d79b5cbca3d 100644 --- a/native/spark-expr/src/csv_funcs/to_csv.rs +++ b/native/spark-expr/src/csv_funcs/to_csv.rs @@ -86,7 +86,7 @@ impl PhysicalExpr for ToCsv { fn evaluate(&self, batch: &RecordBatch) -> Result { let input_array = self.expr.evaluate(batch)?.into_array(batch.num_rows())?; - let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, &self.timezone, false); + let mut cast_options = SparkCastOptions::new(EvalMode::Legacy, &self.timezone); cast_options.null_string = self.csv_write_options.null_value.clone(); let struct_array = as_struct_array(&input_array); diff --git a/native/spark-expr/src/json_funcs/to_json.rs b/native/spark-expr/src/json_funcs/to_json.rs index caf87514fad..94f29e3f458 100644 --- a/native/spark-expr/src/json_funcs/to_json.rs +++ b/native/spark-expr/src/json_funcs/to_json.rs @@ -139,7 +139,7 @@ fn array_to_json_string( spark_cast( ColumnarValue::Array(Arc::clone(arr)), &DataType::Utf8, - &SparkCastOptions::new(EvalMode::Legacy, timezone, false), + &SparkCastOptions::new(EvalMode::Legacy, timezone), )? .into_array(arr.len()) } diff --git a/native/spark-expr/src/math_funcs/internal/decimal_rescale_check.rs b/native/spark-expr/src/math_funcs/internal/decimal_rescale_check.rs index 190d7b060c7..580a01c9afd 100644 --- a/native/spark-expr/src/math_funcs/internal/decimal_rescale_check.rs +++ b/native/spark-expr/src/math_funcs/internal/decimal_rescale_check.rs @@ -195,7 +195,7 @@ impl PhysicalExpr for DecimalRescaleCheckOverflow { Err(_) if self.fail_on_error => spark_cast( arg, &target_type, - &SparkCastOptions::new_without_timezone(EvalMode::Ansi, false), + &SparkCastOptions::new_without_timezone(EvalMode::Ansi), ), Err(error) => Err(error), } diff --git a/native/spark-expr/src/math_funcs/modulo_expr.rs b/native/spark-expr/src/math_funcs/modulo_expr.rs index 5c0652e9a01..09d006c9e49 100644 --- a/native/spark-expr/src/math_funcs/modulo_expr.rs +++ b/native/spark-expr/src/math_funcs/modulo_expr.rs @@ -163,14 +163,14 @@ pub fn create_modulo_expr( let left_256 = Arc::new(Cast::new( left, DataType::Decimal256(p1, s1), - SparkCastOptions::new_without_timezone(EvalMode::Legacy, false), + SparkCastOptions::new_without_timezone(EvalMode::Legacy), None, None, )); let right_256 = Arc::new(Cast::new( right_non_ansi_safe, DataType::Decimal256(p2, s2), - SparkCastOptions::new_without_timezone(EvalMode::Legacy, false), + SparkCastOptions::new_without_timezone(EvalMode::Legacy), None, None, )); @@ -192,7 +192,7 @@ pub fn create_modulo_expr( Ok(Arc::new(Cast::new( modulo_scalar_func, data_type, - SparkCastOptions::new_without_timezone(EvalMode::Legacy, false), + SparkCastOptions::new_without_timezone(EvalMode::Legacy), None, None, ))) diff --git a/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala b/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala index e14fb086b12..b6fcdf32806 100644 --- a/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala +++ b/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala @@ -23,7 +23,6 @@ import org.apache.spark.sql.catalyst.expressions.{Attribute, Cast, Expression, L import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.{ArrayType, DataType, DataTypes, DecimalType, MapType, NullType, StructType, TimestampNTZType, TimestampType} -import org.apache.comet.CometConf import org.apache.comet.CometSparkSessionExtensions.{isSpark40Plus, withFallbackReason} import org.apache.comet.DataTypeSupport.isComplexType import org.apache.comet.serde.{CodegenDispatchFallback, CometExpressionSerde, CometTimeZone, Compatible, ExprOuterClass, Incompatible, SupportLevel, Unsupported} @@ -167,10 +166,6 @@ object CometCast castBuilder.setChild(childExpr) castBuilder.setDatatype(dataType) castBuilder.setEvalMode(evalModeToProto(evalMode)) - castBuilder.setAllowIncompat( - SQLConf.get - .getConfString(CometConf.getExprAllowIncompatConfigKey(classOf[Cast]), "false") - .toBoolean) castBuilder.setTimezone(timeZone) castBuilder.setIsSpark4Plus(isSpark40Plus) Some( diff --git a/spark/src/main/scala/org/apache/comet/expressions/CometEvalMode.scala b/spark/src/main/scala/org/apache/comet/expressions/CometEvalMode.scala index 59e9c89a6b5..d61f0eb3c7b 100644 --- a/spark/src/main/scala/org/apache/comet/expressions/CometEvalMode.scala +++ b/spark/src/main/scala/org/apache/comet/expressions/CometEvalMode.scala @@ -31,12 +31,4 @@ package org.apache.comet.expressions */ object CometEvalMode extends Enumeration { val LEGACY, ANSI, TRY = Value - - def fromBoolean(ansiEnabled: Boolean): Value = if (ansiEnabled) { - ANSI - } else { - LEGACY - } - - def fromString(str: String): CometEvalMode.Value = CometEvalMode.withName(str) } diff --git a/spark/src/main/scala/org/apache/comet/serde/datetime.scala b/spark/src/main/scala/org/apache/comet/serde/datetime.scala index 0317132ef00..8b3ddab5085 100644 --- a/spark/src/main/scala/org/apache/comet/serde/datetime.scala +++ b/spark/src/main/scala/org/apache/comet/serde/datetime.scala @@ -69,7 +69,6 @@ trait CometExprGetDateField[T <: GetDateField] { .setChild(e) .setDatatype(serializeDataType(IntegerType).get) .setEvalMode(ExprOuterClass.EvalMode.LEGACY) - .setAllowIncompat(false) .build()) .build() }) @@ -503,7 +502,6 @@ object CometUnixDate extends CometExpressionSerde[UnixDate] { .setChild(child) .setDatatype(serializeDataType(IntegerType).get) .setEvalMode(ExprOuterClass.EvalMode.LEGACY) - .setAllowIncompat(false) .build()) .build() } @@ -872,7 +870,6 @@ object CometDays extends CometExpressionSerde[Days] { .setChild(dateExpr) .setDatatype(serializeDataType(IntegerType).get) .setEvalMode(ExprOuterClass.EvalMode.LEGACY) - .setAllowIncompat(false) .build()) .build() }