From 73147eb8cdfbea860b61d8fef4e36d50ece9532a Mon Sep 17 00:00:00 2001 From: sam-1112 Date: Wed, 23 Sep 2026 16:50:58 +0800 Subject: [PATCH 1/2] refactor: remove unused cast compatibility flag --- .../expression-audits/conversion_funcs.md | 2 +- native/core/src/execution/planner.rs | 8 +-- native/proto/src/proto/expr.proto | 3 +- .../benches/cast_binary_to_string.rs | 2 +- .../benches/cast_decimal_to_boolean.rs | 2 +- .../benches/cast_decimal_to_string.rs | 2 +- .../benches/cast_float_to_decimal.rs | 2 +- .../benches/cast_float_to_string.rs | 2 +- .../spark-expr/benches/cast_from_boolean.rs | 2 +- native/spark-expr/benches/cast_from_string.rs | 8 +-- .../spark-expr/benches/cast_int_to_decimal.rs | 2 +- .../benches/cast_int_to_timestamp.rs | 2 +- .../benches/cast_non_int_numeric_timestamp.rs | 2 +- native/spark-expr/benches/cast_numeric.rs | 4 +- .../spark-expr/benches/cast_string_to_date.rs | 2 +- .../benches/cast_string_to_timestamp.rs | 4 +- native/spark-expr/benches/to_csv.rs | 2 +- native/spark-expr/benches/wide_decimal.rs | 2 +- .../src/conversion_funcs/boolean.rs | 2 +- .../spark-expr/src/conversion_funcs/cast.rs | 57 +++++++------------ .../src/conversion_funcs/numeric.rs | 2 +- .../spark-expr/src/conversion_funcs/string.rs | 6 +- .../src/conversion_funcs/temporal.rs | 12 ++-- native/spark-expr/src/csv_funcs/to_csv.rs | 2 +- native/spark-expr/src/json_funcs/to_json.rs | 2 +- .../spark-expr/src/math_funcs/modulo_expr.rs | 6 +- .../apache/comet/expressions/CometCast.scala | 4 -- .../comet/expressions/CometEvalMode.scala | 8 --- .../org/apache/comet/serde/datetime.scala | 3 - 29 files changed, 64 insertions(+), 93 deletions(-) diff --git a/docs/source/contributor-guide/expression-audits/conversion_funcs.md b/docs/source/contributor-guide/expression-audits/conversion_funcs.md index ea9a9bf737e..56ac48b58dd 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/src/execution/planner.rs b/native/core/src/execution/planner.rs index 8f030da455b..e8f78f5b91f 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -652,7 +652,6 @@ impl PhysicalPlanner { SparkCastOptions::new_with_version( eval_mode, &expr.timezone, - expr.allow_incompat, expr.is_spark4_plus, ), spark_expr.expr_id, @@ -718,7 +717,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, ))), @@ -838,8 +837,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 = @@ -4702,7 +4700,7 @@ fn create_case_expr( if let Some(coerce_type) = get_coerce_type_for_case_expression(&then_types, else_type.as_ref()) { - let cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy, false); + let cast_options = SparkCastOptions::new_without_timezone(EvalMode::Legacy); let when_then_pairs = when_then_pairs .iter() diff --git a/native/proto/src/proto/expr.proto b/native/proto/src/proto/expr.proto index 34d502daac4..93bebb5d80b 100644 --- a/native/proto/src/proto/expr.proto +++ b/native/proto/src/proto/expr.proto @@ -402,7 +402,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_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 e28afb82725..c1b3b1dc4b5 100644 --- a/native/spark-expr/benches/cast_string_to_timestamp.rs +++ b/native/spark-expr/benches/cast_string_to_timestamp.rs @@ -168,7 +168,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, ); @@ -188,7 +188,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/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 ec932ae68f9..818aa2c3952 100644 --- a/native/spark-expr/benches/wide_decimal.rs +++ b/native/spark-expr/benches/wide_decimal.rs @@ -60,7 +60,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/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 7a9b93ff163..03f96a1ee8d 100644 --- a/native/spark-expr/src/conversion_funcs/cast.rs +++ b/native/spark-expr/src/conversion_funcs/cast.rs @@ -130,8 +130,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, @@ -147,11 +145,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, @@ -160,11 +157,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, @@ -173,15 +169,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) } } } @@ -862,7 +853,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(); @@ -884,7 +875,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(); @@ -908,7 +899,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(); @@ -930,7 +921,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)), @@ -994,7 +985,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, @@ -1013,7 +1004,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, @@ -1025,7 +1016,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, @@ -1051,7 +1042,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::(); @@ -1085,7 +1076,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 { @@ -1123,7 +1114,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()); @@ -1148,7 +1139,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 { @@ -1174,11 +1165,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)); @@ -1197,11 +1186,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)); @@ -1224,7 +1211,7 @@ mod tests { let to_array = cast_array( from_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 81972425e89..bdf0c41ddd7 100644 --- a/native/spark-expr/src/conversion_funcs/numeric.rs +++ b/native/spark-expr/src/conversion_funcs/numeric.rs @@ -2316,7 +2316,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 b466a5f4a09..366a64099b1 100644 --- a/native/spark-expr/src/conversion_funcs/string.rs +++ b/native/spark-expr/src/conversion_funcs/string.rs @@ -2265,7 +2265,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 { @@ -2332,7 +2332,7 @@ mod tests { for to_type in &to_types { let input = format!("{pad}2020-01-01 12:34:56{pad}"); let array: ArrayRef = Arc::new(StringArray::from(vec![Some(input.as_str())])); - let options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false); + let options = SparkCastOptions::new(EvalMode::Legacy, "UTC"); let result = cast_array(array, to_type, &options).unwrap(); assert!( !result.is_null(0), @@ -2527,7 +2527,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/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 d29ef7cd3b7..258adb425d2 100644 --- a/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala +++ b/spark/src/main/scala/org/apache/comet/expressions/CometCast.scala @@ -155,10 +155,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(timeZoneId.getOrElse("UTC")) 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 d1ca21b8869..1fefc1677d4 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() }) @@ -488,7 +487,6 @@ object CometUnixDate extends CometExpressionSerde[UnixDate] { .setChild(child) .setDatatype(serializeDataType(IntegerType).get) .setEvalMode(ExprOuterClass.EvalMode.LEGACY) - .setAllowIncompat(false) .build()) .build() } @@ -851,7 +849,6 @@ object CometDays extends CometExpressionSerde[Days] { .setChild(dateExpr) .setDatatype(serializeDataType(IntegerType).get) .setEvalMode(ExprOuterClass.EvalMode.LEGACY) - .setAllowIncompat(false) .build()) .build() } From 4cedd802ac87b73c50d1b74770e7d4856366485a Mon Sep 17 00:00:00 2001 From: sam-1112 Date: Thu, 24 Sep 2026 02:35:05 +0800 Subject: [PATCH 2/2] refactor: drop unused parquet allow_incompat option The flag has no remaining reader after the cast-option cleanup, and the unused CometConf import fails scalafix. --- .../benches/parquet_timestamp_conversion.rs | 2 +- .../src/execution/operators/iceberg_scan.rs | 2 +- native/core/src/execution/planner.rs | 2 +- native/core/src/parquet/cast_column.rs | 2 +- native/core/src/parquet/parquet_exec.rs | 3 +- native/core/src/parquet/parquet_support.rs | 20 +++++------ native/core/src/parquet/schema_adapter.rs | 33 +++++++++---------- .../apache/comet/expressions/CometCast.scala | 1 - 8 files changed, 29 insertions(+), 36 deletions(-) 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/iceberg_scan.rs b/native/core/src/execution/operators/iceberg_scan.rs index fdcd6c8ba5a..e74e64a5667 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 = diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index e8f78f5b91f..defaf537430 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -6723,7 +6723,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 8c72455f4db..92ed6b1f316 100644 --- a/native/core/src/parquet/cast_column.rs +++ b/native/core/src/parquet/cast_column.rs @@ -385,7 +385,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/parquet_exec.rs b/native/core/src/parquet/parquet_exec.rs index 5d70f2aaa76..4a7c739073b 100644 --- a/native/core/src/parquet/parquet_exec.rs +++ b/native/core/src/parquet/parquet_exec.rs @@ -324,8 +324,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 d6ac9bed636..ffb89537171 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -83,8 +83,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, @@ -138,11 +135,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, @@ -1383,7 +1379,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!( @@ -1491,7 +1487,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), @@ -1611,7 +1607,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)); @@ -1660,7 +1656,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()), @@ -1800,7 +1796,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( @@ -1896,7 +1892,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, diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index e06408b44e1..76fa9d5173d 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -1353,7 +1353,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; @@ -2038,7 +2037,7 @@ 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( @@ -2083,7 +2082,7 @@ 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( @@ -2170,7 +2169,7 @@ 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() @@ -2685,7 +2684,7 @@ 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), @@ -2720,7 +2719,7 @@ 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!() @@ -2841,7 +2840,7 @@ 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 @@ -2868,7 +2867,7 @@ 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 @@ -2894,7 +2893,7 @@ 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 { @@ -2912,7 +2911,7 @@ 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()?; @@ -3019,7 +3018,7 @@ 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( @@ -3079,7 +3078,7 @@ 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)?; @@ -3174,7 +3173,7 @@ 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) @@ -3205,7 +3204,7 @@ 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)) @@ -3249,7 +3248,7 @@ 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) @@ -3319,7 +3318,7 @@ 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, @@ -3430,7 +3429,7 @@ 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 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 258adb425d2..74f3b66b131 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, Compatible, ExprOuterClass, Incompatible, SupportLevel, Unsupported}