Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion native/core/benches/parquet_timestamp_conversion.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion native/core/src/execution/operators/iceberg_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
10 changes: 4 additions & 6 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -653,7 +653,6 @@ impl PhysicalPlanner {
SparkCastOptions::new_with_version(
eval_mode,
&expr.timezone,
expr.allow_incompat,
expr.is_spark4_plus,
),
spark_expr.expr_id,
Expand Down Expand Up @@ -719,7 +718,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,
))),
Expand Down Expand Up @@ -839,8 +838,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 =
Expand Down Expand Up @@ -4710,7 +4708,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()
Expand Down Expand Up @@ -6733,7 +6731,7 @@ mod tests {
.with_table_parquet_options(TableParquetOptions::new()),
) as Arc<dyn FileSource>;

let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false);
let spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC");

let expr_adapter_factory: Arc<dyn PhysicalExprAdapterFactory> = Arc::new(
SparkPhysicalExprAdapterFactory::new(spark_parquet_options, None),
Expand Down
2 changes: 1 addition & 1 deletion native/core/src/parquet/cast_column.rs
Original file line number Diff line number Diff line change
Expand Up @@ -385,7 +385,7 @@ mod tests {
let expr: Arc<dyn PhysicalExpr> = 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());
Expand Down
3 changes: 1 addition & 2 deletions native/core/src/parquet/parquet_exec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
20 changes: 8 additions & 12 deletions native/core/src/parquet/parquet_support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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> -> int must be an error");
assert!(
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -1660,7 +1656,7 @@ mod tests {

for timezone in [None::<Arc<str>>, 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()),
Expand Down Expand Up @@ -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::<Arc<str>>, Some(Arc::from("UTC"))] {
for overflow in [i64::MIN, i64::MAX] {
let values: ArrayRef = Arc::new(
Expand Down Expand Up @@ -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,
Expand Down
33 changes: 16 additions & 17 deletions native/core/src/parquet/schema_adapter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<dyn PhysicalExprAdapterFactory> = Arc::new(
Expand Down Expand Up @@ -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<dyn PhysicalExprAdapterFactory> = Arc::new(
Expand Down Expand Up @@ -2170,7 +2169,7 @@ mod test {
batch: &RecordBatch,
required_schema: SchemaRef,
) -> Result<RecordBatch, DataFusionError> {
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()
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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!()
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -2894,7 +2893,7 @@ mod test {
Arc::new(Int32Array::from(Vec::<i32>::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 {
Expand All @@ -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()?;
Expand Down Expand Up @@ -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<dyn PhysicalExprAdapterFactory> = Arc::new(
Expand Down Expand Up @@ -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)?;
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion native/proto/src/proto/expr.proto
Original file line number Diff line number Diff line change
Expand Up @@ -417,7 +417,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;
Expand Down
2 changes: 1 addition & 1 deletion native/spark-expr/benches/cast_binary_to_string.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ fn build(size: usize, width: usize, null_every: usize) -> ArrayRef {
}

fn options(style: Option<BinaryOutputStyle>) -> SparkCastOptions {
let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false);
let mut options = SparkCastOptions::new(EvalMode::Legacy, "UTC");
options.binary_output_style = style;
options
}
Expand Down
Loading