From 7594f111a214f619c148b98df5b0bc367f75c845 Mon Sep 17 00:00:00 2001 From: Nat Wilson Date: Fri, 2 Oct 2026 09:34:55 -0700 Subject: [PATCH] improve Spark from_utc_timestamp compatibility The timezone spellings accepted by Spark are different than those accepted by Arrow's parser. As a result, `from_utc_timestamp` in comet isn't safe to use, since it would fail on things like Java short IDs, prefixed offsets with second precision, and SystemV-style IDs. Additionally, Arrow will accept offsets beyond the Spark limit of 18 hours. This change improves compatibility by parsing timezone IDs with the rules implemented by Spark. As part of doing this, I added SparkFromUtcTimestampExpr so that computed timezones evaluated for rows with non-null timestamps, avoiding failing when a timezone is invalid but wouldn't be required anyway. (Full compatibility with spark behaviour seems quite complex and I haven't attempted it here.) It also adds more documentation recording remaining differences: - timezone validation and evaluation order can differ from Spark - chrono accepts a narrower calendar range than Spark implemented range - chrono-tz doesn't currently perform DST transitions past 2099 This does not achieve full compatibility with Spark due to those issues, but it makes the region of incompatibility a lot smaller. --- .../functions/src/datetime/to_local_time.rs | 28 +- datafusion/spark/Cargo.toml | 4 + .../spark/benches/from_utc_timestamp.rs | 94 +++ .../function/datetime/from_utc_timestamp.rs | 576 +++++++++++++++++- datafusion/spark/src/function/datetime/mod.rs | 1 + .../spark/src/function/datetime/timezone.rs | 442 ++++++++++++++ .../spark/datetime/from_utc_timestamp.slt | 25 + typos.toml | 1 + 8 files changed, 1149 insertions(+), 22 deletions(-) create mode 100644 datafusion/spark/benches/from_utc_timestamp.rs create mode 100644 datafusion/spark/src/function/datetime/timezone.rs diff --git a/datafusion/functions/src/datetime/to_local_time.rs b/datafusion/functions/src/datetime/to_local_time.rs index 973afa549e0b2..1e98f2d654033 100644 --- a/datafusion/functions/src/datetime/to_local_time.rs +++ b/datafusion/functions/src/datetime/to_local_time.rs @@ -15,7 +15,6 @@ // specific language governing permissions and limitations // under the License. -use std::ops::Add; use std::sync::Arc; use arrow::array::timezone::Tz; @@ -30,8 +29,8 @@ use chrono::{DateTime, MappedLocalTime, Offset, TimeDelta, TimeZone, Utc}; use datafusion_common::cast::as_primitive_array; use datafusion_common::{ - Result, ScalarValue, exec_err, internal_datafusion_err, internal_err, - utils::take_function_args, + Result, ScalarValue, exec_datafusion_err, exec_err, internal_datafusion_err, + internal_err, utils::take_function_args, }; use datafusion_expr::{ Coercion, ColumnarValue, Documentation, ScalarFunctionArgs, ScalarUDFImpl, Signature, @@ -307,7 +306,10 @@ fn to_local_time(time_value: &ColumnarValue) -> Result { /// ``` /// /// See `test_adjust_to_local_time()` for example -pub fn adjust_to_local_time(ts: i64, tz: Tz) -> Result { +pub fn adjust_to_local_time( + ts: i64, + tz: impl TimeZone + Copy, +) -> Result { fn convert_timestamp(ts: i64, converter: F) -> Result> where F: Fn(i64) -> MappedLocalTime>, @@ -337,12 +339,18 @@ pub fn adjust_to_local_time(ts: i64, tz: Tz) -> Result RecordBatch { + let mut fields = vec![ + Field::new( + "timestamp", + DataType::Timestamp(TimeUnit::Microsecond, None), + true, + ), + Field::new("timezone", DataType::Utf8, false), + ]; + let mut arrays: Vec = vec![ + Arc::new(TimestampMicrosecondArray::from_iter((0..rows).map(|row| { + (!half_null || row % 2 != 0).then_some(row as i64 * 86_400_000_000) + }))), + Arc::new(StringArray::from_iter_values( + (0..rows).map(|row| if (row / 2) % 2 == 0 { "PST" } else { "EST" }), + )), + ]; + + for column in 2..columns { + fields.push(Field::new( + format!("unused_{column}"), + DataType::Int64, + false, + )); + arrays.push(Arc::new(Int64Array::from_iter_values( + (0..rows).map(|row| (column * rows + row) as i64), + ))); + } + + RecordBatch::try_new(Arc::new(Schema::new(fields)), arrays).unwrap() +} + +fn criterion_benchmark(c: &mut Criterion) { + let rows = 8192; + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Column::new("timestamp", 0)), + Arc::new(Column::new("timezone", 1)), + ); + let mut group = c.benchmark_group("from_utc_timestamp_column"); + group.throughput(Throughput::Elements(rows as u64)); + + for columns in [2, 64, 256] { + for (half_null, name) in [(false, "no_nulls"), (true, "half_nulls")] { + let batch = batch(rows, columns, half_null); + let result = expr.evaluate(&batch).unwrap().into_array(rows).unwrap(); + assert_eq!(result.len(), rows); + assert_eq!(result.null_count(), if half_null { rows / 2 } else { 0 }); + + group.bench_with_input( + BenchmarkId::new(format!("{columns}_columns"), name), + &batch, + |b, batch| { + b.iter(|| black_box(expr.evaluate(black_box(batch)).unwrap())); + }, + ); + } + } + group.finish(); +} + +criterion_group!(benches, criterion_benchmark); +criterion_main!(benches); diff --git a/datafusion/spark/src/function/datetime/from_utc_timestamp.rs b/datafusion/spark/src/function/datetime/from_utc_timestamp.rs index bfca677c1dcce..2e18be931fac4 100644 --- a/datafusion/spark/src/function/datetime/from_utc_timestamp.rs +++ b/datafusion/spark/src/function/datetime/from_utc_timestamp.rs @@ -15,24 +15,30 @@ // specific language governing permissions and limitations // under the License. +use std::fmt::{Display, Formatter}; +use std::hash::{Hash, Hasher}; use std::sync::Arc; -use arrow::array::timezone::Tz; -use arrow::array::{Array, ArrayRef, AsArray, PrimitiveBuilder, StringArrayType}; +use super::timezone::SparkTimeZone; +use arrow::array::{ + Array, ArrayRef, AsArray, BooleanArray, PrimitiveBuilder, StringArrayType, + new_empty_array, +}; use arrow::datatypes::TimeUnit; use arrow::datatypes::{ - ArrowTimestampType, DataType, Field, FieldRef, TimestampMicrosecondType, + ArrowTimestampType, DataType, Field, FieldRef, Schema, TimestampMicrosecondType, TimestampMillisecondType, TimestampNanosecondType, TimestampSecondType, }; +use arrow::record_batch::RecordBatch; use datafusion_common::types::{NativeType, logical_string}; use datafusion_common::utils::take_function_args; use datafusion_common::{Result, exec_datafusion_err, exec_err, internal_err}; use datafusion_expr::{ - Coercion, ColumnarValue, ReturnFieldArgs, ScalarFunctionArgs, ScalarUDFImpl, - Signature, TypeSignatureClass, Volatility, + Coercion, ColumnarValue, ExpressionPlacement, ReturnFieldArgs, ScalarFunctionArgs, + ScalarUDFImpl, Signature, TypeSignatureClass, Volatility, }; -use datafusion_functions::datetime::to_local_time::adjust_to_local_time; use datafusion_functions::utils::make_scalar_function; +use datafusion_physical_expr_common::physical_expr::PhysicalExpr; /// Apache Spark `from_utc_timestamp` function. /// @@ -42,6 +48,21 @@ use datafusion_functions::utils::make_scalar_function; /// timezone-agnostic. So in Apache Spark this function just shift the timestamp value from UTC timezone to /// the given timezone. /// +/// # Compatibility +/// +/// This function accepts Spark's additional offset spellings, Java short IDs, and +/// legacy SystemV zone IDs. They are parsed only when both arguments are non-null. +/// Spark's generated execution may validate a foldable timezone before processing +/// rows, so this implementation can return nulls where Spark raises an error. +/// +/// Timezone conversion uses Chrono's calendar range, which is narrower than Spark's +/// implemented timestamp range, although it includes Spark's documented range. +/// Region timezone rules from `chrono-tz` may also differ after 2099 because it does +/// not extrapolate later daylight saving time transitions. +/// +/// This UDF operates on already-evaluated arguments. Callers that need conditional +/// evaluation of computed timezone expressions should construct [`SparkFromUtcTimestampExpr`]. +/// /// See #[derive(Debug, PartialEq, Eq, Hash)] pub struct SparkFromUtcTimestamp { @@ -169,16 +190,21 @@ where { let ts_primitive = ts_array.as_primitive::(); let mut builder = PrimitiveBuilder::::with_capacity(ts_array.len()); + // Parse a scalar zone once per batch and avoid re-parsing consecutive column values + let mut last_timezone = None; for (ts_opt, tz_opt) in ts_primitive.iter().zip(tz_array.iter()) { match (ts_opt, tz_opt) { (Some(ts), Some(tz_str)) => { - let tz: Tz = tz_str.parse().map_err(|e| { - exec_datafusion_err!( - "`from_utc_timestamp`: invalid timezone '{tz_str}': {e}" - ) - })?; - let val = adjust_to_local_time::(ts, tz)?; + let timezone = match last_timezone { + Some((previous, timezone)) if previous == tz_str => timezone, + _ => { + let timezone = parse_timezone(tz_str)?; + last_timezone = Some((tz_str, timezone)); + timezone + } + }; + let val = timezone.adjust_to_local_time::(ts)?; builder.append_value(val); } _ => builder.append_null(), @@ -188,3 +214,529 @@ where builder = builder.with_timezone_opt(return_tz_opt); Ok(Arc::new(builder.finish())) } + +fn parse_timezone(timezone: &str) -> Result { + SparkTimeZone::parse(timezone).map_err(|e| { + exec_datafusion_err!("`from_utc_timestamp`: invalid timezone '{timezone}': {e}") + }) +} + +/// Physical `from_utc_timestamp` expression with conditional timezone evaluation. +/// +/// The timestamp child is evaluated once for each non-empty batch. If every timestamp is null, +/// the timezone child is skipped. Otherwise, literal and column timezones are read directly, and +/// their values are ignored on null timestamp rows. Computed timezone expressions are evaluated +/// only for rows with non-null timestamps. Empty batches evaluate neither child. +/// +/// This intentionally differs from Spark's generated evaluation of foldable timezones, which +/// evaluates the timezone and validates a non-null value before processing timestamp rows. This +/// expression can therefore skip an invalid timezone or an error in the timezone expression when +/// no timestamp is non-null. It also evaluates the timestamp first when the timezone is null, so a +/// timestamp error can be raised even though Spark's generated path would return null. Planners +/// using this expression should expose these differences as an incompatibility. +#[derive(Debug, Clone)] +pub struct SparkFromUtcTimestampExpr { + timestamp: Arc, + timezone: Arc, +} + +impl SparkFromUtcTimestampExpr { + /// Construct an expression that skips computed timezone evaluation for null timestamps. + pub fn new( + timestamp: Arc, + timezone: Arc, + ) -> Self { + Self { + timestamp, + timezone, + } + } +} + +impl Display for SparkFromUtcTimestampExpr { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!( + f, + "from_utc_timestamp({}, {})", + self.timestamp, self.timezone + ) + } +} + +impl PartialEq for SparkFromUtcTimestampExpr { + fn eq(&self, other: &Self) -> bool { + self.timestamp.eq(&other.timestamp) && self.timezone.eq(&other.timezone) + } +} + +impl Eq for SparkFromUtcTimestampExpr {} + +impl Hash for SparkFromUtcTimestampExpr { + fn hash(&self, state: &mut H) { + self.timestamp.hash(state); + self.timezone.hash(state); + } +} + +impl PhysicalExpr for SparkFromUtcTimestampExpr { + fn fmt_sql(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + Display::fmt(self, f) + } + + fn return_field(&self, input_schema: &Schema) -> Result { + let timestamp = self.timestamp.return_field(input_schema)?; + Ok(Arc::new(Field::new( + "from_utc_timestamp", + timestamp.data_type().clone(), + timestamp.is_nullable() || self.timezone.nullable(input_schema)?, + ))) + } + + fn evaluate(&self, batch: &RecordBatch) -> Result { + // There are no rows on which either child should run. + if batch.num_rows() == 0 { + return Ok(ColumnarValue::Array(new_empty_array( + &self.data_type(batch.schema_ref())?, + ))); + } + + let timestamp = self.timestamp.evaluate(batch)?; + let timezone = match ×tamp { + ColumnarValue::Scalar(value) if value.is_null() => { + return Ok(timestamp); + } + ColumnarValue::Scalar(_) => self.timezone.evaluate(batch)?, + ColumnarValue::Array(array) if array.null_count() == array.len() => { + return Ok(timestamp); + } + ColumnarValue::Array(_) + if matches!( + self.timezone.placement(), + ExpressionPlacement::Literal | ExpressionPlacement::Column + ) => + { + // Literals and columns only read existing values. The conversion loop skips null + // timestamps, so avoid filtering and copying unrelated batch columns. + self.timezone.evaluate(batch)? + } + ColumnarValue::Array(array) => match array.nulls() { + Some(nulls) => { + let selection = BooleanArray::new(nulls.inner().clone(), None); + self.timezone.evaluate_selection(batch, &selection)? + } + None => self.timezone.evaluate(batch)?, + }, + }; + make_scalar_function(spark_from_utc_timestamp, vec![])(&[timestamp, timezone]) + } + + fn children(&self) -> Vec<&Arc> { + vec![&self.timestamp, &self.timezone] + } + + fn with_new_children( + self: Arc, + children: Vec>, + ) -> Result> { + let [timestamp, timezone] = children.as_slice() else { + return internal_err!( + "SparkFromUtcTimestampExpr expected 2 children, got {}", + children.len() + ); + }; + Ok(Arc::new(Self::new( + Arc::clone(timestamp), + Arc::clone(timezone), + ))) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow::array::{ + LargeStringArray, StringArray, StringViewArray, TimestampMicrosecondArray, + }; + use arrow::compute::CastOptions; + use chrono::{DateTime, Utc}; + use datafusion::physical_expr::expressions::{ + BinaryExpr, CaseExpr, CastExpr, Column, Literal, + }; + use datafusion_common::{ScalarValue, config::ConfigOptions}; + use datafusion_expr::Operator; + + fn expression_batch(timestamps: Vec>, values: Vec<&str>) -> RecordBatch { + RecordBatch::try_new( + Arc::new(Schema::new(vec![ + Field::new( + "timestamp", + DataType::Timestamp(TimeUnit::Microsecond, Some(Arc::from("UTC"))), + true, + ), + Field::new("value", DataType::Utf8, false), + ])), + vec![ + Arc::new( + TimestampMicrosecondArray::from(timestamps).with_timezone("UTC"), + ), + Arc::new(StringArray::from(values)), + ], + ) + .unwrap() + } + + fn throwing_timestamp() -> Arc { + Arc::new(CastExpr::new( + Arc::new(Literal::new(ScalarValue::Utf8(Some("invalid".to_owned())))), + DataType::Timestamp(TimeUnit::Microsecond, Some(Arc::from("UTC"))), + Some(CastOptions { + safe: false, + ..Default::default() + }), + )) + } + + fn computed_timezone() -> Arc { + let value = Arc::new(CastExpr::new( + Arc::new(Column::new("value", 1)), + DataType::Int32, + Some(CastOptions { + safe: false, + ..Default::default() + }), + )); + let is_zero = Arc::new(BinaryExpr::new( + value, + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(0)))), + )); + Arc::new( + CaseExpr::try_new( + None, + vec![( + is_zero, + Arc::new(Literal::new(ScalarValue::Utf8(Some("UTC".to_owned())))), + )], + Some(Arc::new(Literal::new(ScalarValue::Utf8(Some( + "PST".to_owned(), + ))))), + ) + .unwrap(), + ) + } + + fn invoke( + timestamp: ColumnarValue, + timezone: ColumnarValue, + ) -> Result { + let return_field = Arc::new(Field::new("result", timestamp.data_type(), true)); + SparkFromUtcTimestamp::new().invoke_with_args(ScalarFunctionArgs { + args: vec![timestamp, timezone], + arg_fields: vec![], + number_rows: 0, + return_field, + config_options: Arc::new(ConfigOptions::default()), + }) + } + + fn timestamps(values: Vec>) -> ColumnarValue { + ColumnarValue::Array(Arc::new( + TimestampMicrosecondArray::from(values).with_timezone("UTC"), + )) + } + + fn scalar_zone(value: Option<&str>) -> ColumnarValue { + ColumnarValue::Scalar(ScalarValue::Utf8(value.map(str::to_owned))) + } + + fn assert_timestamps(result: ColumnarValue, expected: &[Option]) { + let ColumnarValue::Array(array) = result else { + panic!("expected timestamp array"); + }; + let array = array.as_primitive::(); + assert_eq!(array.iter().collect::>(), expected); + assert_eq!(array.timezone(), Some("UTC")); + } + + #[test] + fn column_zones_preserve_values_nulls_and_timestamp_metadata() { + let summer = 1_719_792_000_123_456; + let zones = vec![ + Some("+00:00:01"), + Some("GMT+1"), + Some("PST"), + Some("invalid"), + None, + ]; + let zone_arrays: Vec = vec![ + Arc::new(StringArray::from(zones.clone())), + Arc::new(LargeStringArray::from(zones.clone())), + Arc::new(StringViewArray::from(zones)), + ]; + for zones in zone_arrays { + let result = invoke( + timestamps(vec![Some(-1), Some(0), Some(summer), None, Some(summer)]), + ColumnarValue::Array(zones), + ) + .unwrap(); + assert_timestamps( + result, + &[ + Some(999_999), + Some(3_600_000_000), + Some(summer - 25_200_000_000), + None, + None, + ], + ); + } + } + + #[test] + fn scalar_and_array_arguments_work_across_batches() { + for values in [vec![Some(-1), None], vec![Some(123_456), Some(-1_000_001)]] { + let expected = values + .iter() + .map(|value| value.map(|value| value + 3_723_000_000)) + .collect::>(); + let result = + invoke(timestamps(values), scalar_zone(Some("UTC+1:02:03"))).unwrap(); + assert_timestamps(result, &expected); + } + + let timestamp = 1_719_792_000_123_456; + let result = invoke( + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond( + Some(timestamp), + Some(Arc::from("UTC")), + )), + ColumnarValue::Array(Arc::new(StringArray::from(vec![ + Some("GMT+1"), + Some("PST"), + None, + Some("GMT+1"), + ]))), + ) + .unwrap(); + assert_timestamps( + result, + &[ + Some(timestamp + 3_600_000_000), + Some(timestamp - 25_200_000_000), + None, + Some(timestamp + 3_600_000_000), + ], + ); + } + + #[test] + fn fixed_offset_seconds_preserve_each_timestamp_unit() { + for (input, expected) in [ + ( + ScalarValue::TimestampSecond(Some(-1), None), + ScalarValue::TimestampSecond(Some(0), None), + ), + ( + ScalarValue::TimestampMillisecond(Some(-1), None), + ScalarValue::TimestampMillisecond(Some(999), None), + ), + ( + ScalarValue::TimestampMicrosecond(Some(-1), None), + ScalarValue::TimestampMicrosecond(Some(999_999), None), + ), + ( + ScalarValue::TimestampNanosecond(Some(-1), None), + ScalarValue::TimestampNanosecond(Some(999_999_999), None), + ), + ] { + let result = invoke( + ColumnarValue::Scalar(input.clone()), + scalar_zone(Some("+00:00:01")), + ) + .unwrap(); + let ColumnarValue::Scalar(actual) = result else { + panic!("expected scalar timestamp"); + }; + assert_eq!(actual, expected); + + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Literal::new(input)), + Arc::new(Literal::new(ScalarValue::Utf8(Some( + "+00:00:01".to_owned(), + )))), + ); + let result = expr + .evaluate(&expression_batch(vec![Some(0)], vec!["0"])) + .unwrap(); + let ColumnarValue::Scalar(actual) = result else { + panic!("expected scalar timestamp"); + }; + assert_eq!(actual, expected); + } + } + + #[test] + fn invalid_literal_timezone_is_validated_lazily() { + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Column::new("timestamp", 0)), + Arc::new(Literal::new(ScalarValue::Utf8(Some("+19:00".to_owned())))), + ); + assert_timestamps( + expr.evaluate(&expression_batch(vec![None, None], vec!["0", "0"])) + .unwrap(), + &[None, None], + ); + let error = expr + .evaluate(&expression_batch(vec![Some(0)], vec!["0"])) + .unwrap_err(); + assert!(error.to_string().contains("invalid timezone '+19:00'")); + } + + #[test] + fn null_timezone_does_not_skip_timestamp_evaluation() { + let expr = SparkFromUtcTimestampExpr::new( + throwing_timestamp(), + Arc::new(Literal::new(ScalarValue::Utf8(None))), + ); + let error = expr + .evaluate(&expression_batch(vec![Some(0)], vec!["0"])) + .unwrap_err(); + assert!(error.to_string().contains("invalid")); + } + + #[test] + fn timezone_evaluates_only_non_null_timestamp_rows() { + let summer = 1_719_792_000_123_456; + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Column::new("timestamp", 0)), + computed_timezone(), + ); + let batch = expression_batch( + vec![None, Some(0), None, Some(summer)], + vec!["invalid", "0", "invalid", "1"], + ); + assert_timestamps( + expr.evaluate(&batch).unwrap(), + &[None, Some(0), None, Some(summer - 25_200_000_000)], + ); + + // The same timezone child must still throw when its timestamp is non-null. + let error = expr + .evaluate(&expression_batch(vec![Some(0)], vec!["invalid"])) + .unwrap_err(); + assert!(error.to_string().contains("invalid")); + } + + #[test] + fn timezone_column_ignores_values_on_null_timestamp_rows() { + let summer = 1_719_792_000_123_456; + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Column::new("timestamp", 0)), + Arc::new(Column::new("value", 1)), + ); + let batch = expression_batch( + vec![None, Some(0), None, Some(summer)], + vec!["invalid", "UTC", "+19:00", "PST"], + ); + assert_timestamps( + expr.evaluate(&batch).unwrap(), + &[None, Some(0), None, Some(summer - 25_200_000_000)], + ); + + let error = expr + .evaluate(&expression_batch(vec![Some(0)], vec!["+19:00"])) + .unwrap_err(); + assert!(error.to_string().contains("invalid timezone '+19:00'")); + } + + #[test] + fn timezone_skips_all_null_and_scalar_null_timestamps() { + let batch = expression_batch(vec![None, None], vec!["invalid", "invalid"]); + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Column::new("timestamp", 0)), + computed_timezone(), + ); + assert_timestamps(expr.evaluate(&batch).unwrap(), &[None, None]); + + let expr = SparkFromUtcTimestampExpr::new( + Arc::new(Literal::new(ScalarValue::TimestampMicrosecond(None, None))), + computed_timezone(), + ); + assert!(matches!( + expr.evaluate(&batch).unwrap(), + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond(None, None)) + )); + } + + #[test] + fn empty_batch_skips_children() { + let batch = expression_batch(vec![], vec![]); + let expr = + SparkFromUtcTimestampExpr::new(throwing_timestamp(), computed_timezone()); + assert_timestamps(expr.evaluate(&batch).unwrap(), &[]); + } + + #[test] + fn fixed_offset_outside_calendar_range_returns_error() { + for (timestamp, timezone) in [ + (DateTime::::MAX_UTC.timestamp_micros(), "+00:00:01"), + (DateTime::::MIN_UTC.timestamp_micros(), "-00:00:01"), + ] { + let error = invoke( + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond( + Some(timestamp), + None, + )), + scalar_zone(Some(timezone)), + ) + .unwrap_err(); + assert!( + error + .to_string() + .contains("exceeds Chrono's calendar range") + ); + } + } + + #[test] + fn invalid_zones_are_parsed_only_for_non_null_pairs() { + let zones = + ColumnarValue::Array(Arc::new(StringArray::from(vec!["invalid", "+19:00"]))); + assert_timestamps( + invoke(timestamps(vec![None, None]), zones).unwrap(), + &[None, None], + ); + + let error = invoke(timestamps(vec![None, Some(0)]), scalar_zone(Some("+19:00"))) + .unwrap_err(); + assert!(error.to_string().contains("invalid timezone '+19:00'")); + + let result = invoke( + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond(None, None)), + scalar_zone(Some("invalid")), + ) + .unwrap(); + assert!(matches!( + result, + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond(None, None)) + )); + + let result = invoke( + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond( + Some(i64::MAX), + None, + )), + scalar_zone(None), + ) + .unwrap(); + assert!(matches!( + result, + ColumnarValue::Scalar(ScalarValue::TimestampMicrosecond(None, None)) + )); + + assert_timestamps( + invoke(timestamps(vec![]), scalar_zone(Some("invalid"))).unwrap(), + &[], + ); + } +} diff --git a/datafusion/spark/src/function/datetime/mod.rs b/datafusion/spark/src/function/datetime/mod.rs index 70ab024c329aa..a365fbe4a971a 100644 --- a/datafusion/spark/src/function/datetime/mod.rs +++ b/datafusion/spark/src/function/datetime/mod.rs @@ -29,6 +29,7 @@ pub mod make_interval; pub mod monthname; pub mod next_day; pub mod time_trunc; +mod timezone; pub mod to_utc_timestamp; pub mod trunc; pub mod unix; diff --git a/datafusion/spark/src/function/datetime/timezone.rs b/datafusion/spark/src/function/datetime/timezone.rs new file mode 100644 index 0000000000000..17971c85ea4ef --- /dev/null +++ b/datafusion/spark/src/function/datetime/timezone.rs @@ -0,0 +1,442 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use arrow::array::timezone::Tz; +use arrow::datatypes::{ArrowTimestampType, TimeUnit}; +use chrono::{DateTime, Datelike, FixedOffset, NaiveDate, Utc}; +use datafusion_common::{Result, exec_datafusion_err, internal_datafusion_err}; +use datafusion_functions::datetime::to_local_time::adjust_to_local_time; + +/// A timezone resolved using Spark's `getZoneId` spelling rules. +/// +/// Fixed offsets are kept as seconds to preserve Spark's second-resolution offsets. +/// Region IDs use Arrow's Chrono-backed timezone database. +/// Java's legacy SystemV zones require their own rules because that database omits them. +#[derive(Debug, Clone, Copy)] +pub(super) enum SparkTimeZone { + FixedOffset(i32), + Region(Tz), + SystemV(i32), +} + +impl SparkTimeZone { + pub(super) fn parse(timezone: &str) -> Result { + // These are the exact aliases in java.time.ZoneId.SHORT_IDS. The regional + // aliases retain their daylight saving rules; EST, HST, and MST are fixed. + let timezone = match timezone { + "Z" => return Ok(Self::FixedOffset(0)), + // The JDK explicitly excludes this legacy IANA alias. + "ROC" => return Err(exec_datafusion_err!("Unknown timezone '{timezone}'")), + "EST" => return Ok(Self::FixedOffset(-5 * 3600)), + "HST" => return Ok(Self::FixedOffset(-10 * 3600)), + "MST" => return Ok(Self::FixedOffset(-7 * 3600)), + "ACT" => "Australia/Darwin", + "AET" => "Australia/Sydney", + "AGT" => "America/Argentina/Buenos_Aires", + "ART" => "Africa/Cairo", + "AST" => "America/Anchorage", + "BET" => "America/Sao_Paulo", + "BST" => "Asia/Dhaka", + "CAT" => "Africa/Harare", + "CNT" => "America/St_Johns", + "CST" => "America/Chicago", + "CTT" => "Asia/Shanghai", + "EAT" => "Africa/Addis_Ababa", + "ECT" => "Europe/Paris", + "IET" => "America/Indiana/Indianapolis", + "IST" => "Asia/Kolkata", + "JST" => "Asia/Tokyo", + "MIT" => "Pacific/Apia", + "NET" => "Asia/Yerevan", + "NST" => "Pacific/Auckland", + "PLT" => "Asia/Karachi", + "PNT" => "America/Phoenix", + "PRT" => "America/Puerto_Rico", + "PST" => "America/Los_Angeles", + "SST" => "Pacific/Guadalcanal", + "VST" => "Asia/Ho_Chi_Minh", + "SystemV/AST4" => return Ok(Self::FixedOffset(-4 * 3600)), + "SystemV/EST5" => return Ok(Self::FixedOffset(-5 * 3600)), + "SystemV/CST6" => return Ok(Self::FixedOffset(-6 * 3600)), + "SystemV/MST7" => return Ok(Self::FixedOffset(-7 * 3600)), + "SystemV/PST8" => return Ok(Self::FixedOffset(-8 * 3600)), + "SystemV/YST9" => return Ok(Self::FixedOffset(-9 * 3600)), + "SystemV/HST10" => return Ok(Self::FixedOffset(-10 * 3600)), + "SystemV/AST4ADT" => return Ok(Self::SystemV(-4 * 3600)), + "SystemV/EST5EDT" => return Ok(Self::SystemV(-5 * 3600)), + "SystemV/CST6CDT" => return Ok(Self::SystemV(-6 * 3600)), + "SystemV/MST7MDT" => return Ok(Self::SystemV(-7 * 3600)), + "SystemV/PST8PDT" => return Ok(Self::SystemV(-8 * 3600)), + "SystemV/YST9YDT" => return Ok(Self::SystemV(-9 * 3600)), + _ => timezone, + }; + + let offset = if timezone.starts_with(['+', '-']) { + Some(timezone) + } else { + ["UTC", "GMT", "UT"].into_iter().find_map(|prefix| { + timezone + .strip_prefix(prefix) + .filter(|suffix| suffix.is_empty() || suffix.starts_with(['+', '-'])) + }) + }; + if let Some(offset) = offset { + if offset.is_empty() { + return Ok(Self::FixedOffset(0)); + } + return parse_offset(offset).map(Self::FixedOffset).ok_or_else(|| { + exec_datafusion_err!("Invalid timezone offset '{offset}'") + }); + } + + Ok(Self::Region(timezone.parse()?)) + } + + pub(super) fn adjust_to_local_time( + self, + ts: i64, + ) -> Result { + match self { + Self::Region(timezone) => adjust_to_local_time::(ts, timezone), + Self::FixedOffset(seconds) => adjust_by_offset::(ts, seconds), + Self::SystemV(standard_offset) => { + adjust_by_offset::(ts, system_v_offset::(ts, standard_offset)?) + } + } + } +} + +/// Parses Java's signed offsets after Spark's single-digit hour/minute normalization. +/// +/// Spark permits one-digit hours before a colon, and one-digit minutes only when +/// minutes are the final component. Java also accepts compact HHMM and HHMMSS forms. +fn parse_offset(offset: &str) -> Option { + let sign = match offset.as_bytes().first()? { + b'+' => 1, + b'-' => -1, + _ => return None, + }; + let value = &offset[1..]; + let (hours, minutes, seconds) = if let Some((hours, rest)) = value.split_once(':') { + let hours = parse_digits(hours, 1, 2)?; + if let Some((minutes, seconds)) = rest.split_once(':') { + ( + hours, + parse_digits(minutes, 2, 2)?, + parse_digits(seconds, 2, 2)?, + ) + } else { + (hours, parse_digits(rest, 1, 2)?, 0) + } + } else { + // Check ASCII before byte slicing so malformed Unicode inputs cannot panic. + if !value.bytes().all(|byte| byte.is_ascii_digit()) { + return None; + } + match value.len() { + 1 | 2 => (parse_digits(value, 1, 2)?, 0, 0), + 4 => ( + parse_digits(&value[..2], 2, 2)?, + parse_digits(&value[2..], 2, 2)?, + 0, + ), + 6 => ( + parse_digits(&value[..2], 2, 2)?, + parse_digits(&value[2..4], 2, 2)?, + parse_digits(&value[4..], 2, 2)?, + ), + _ => return None, + } + }; + if hours > 18 + || minutes > 59 + || seconds > 59 + || (hours == 18 && (minutes != 0 || seconds != 0)) + { + return None; + } + Some(sign * (hours * 3600 + minutes * 60 + seconds)) +} + +fn parse_digits(value: &str, min_len: usize, max_len: usize) -> Option { + if !(min_len..=max_len).contains(&value.len()) + || !value.bytes().all(|byte| byte.is_ascii_digit()) + { + return None; + } + value.parse().ok() +} + +fn units_per_second() -> i64 { + match T::UNIT { + TimeUnit::Second => 1, + TimeUnit::Millisecond => 1_000, + TimeUnit::Microsecond => 1_000_000, + TimeUnit::Nanosecond => 1_000_000_000, + } +} + +fn adjust_by_offset(ts: i64, offset_seconds: i32) -> Result { + let timezone = FixedOffset::east_opt(offset_seconds).ok_or_else(|| { + internal_datafusion_err!("Spark timezone offsets must be within 18 hours") + })?; + adjust_to_local_time::(ts, timezone) +} + +/// Java's SystemV daylight zones use the same calendar with different standard offsets. +/// The JDK starts these transitions in 1900. Most years change on the last Sundays of +/// April and October; 1974 and 1975 have exceptional start/end dates. +/// The rule definitions are in OpenJDK's `make/data/tzdata/jdk11_backward`. +fn system_v_offset(ts: i64, standard_offset: i32) -> Result { + let seconds = ts.div_euclid(units_per_second::()); + let year = DateTime::::from_timestamp(seconds, 0) + .ok_or_else(|| { + exec_datafusion_err!("Timestamp is outside Chrono's calendar range") + })? + .year(); + if year < 1900 { + return Ok(standard_offset); + } + + let start = match year { + 1974 => NaiveDate::from_ymd_opt(year, 1, 6), + 1975 => NaiveDate::from_ymd_opt(year, 2, 23), + _ => last_sunday(year, 4, 30), + }; + let end = if year == 1974 { + last_sunday(year, 11, 30) + } else { + last_sunday(year, 10, 31) + }; + let transition = |date: Option, offset_before: i32| { + date.and_then(|date| date.and_hms_opt(2, 0, 0)) + .map(|local| local.and_utc().timestamp() - i64::from(offset_before)) + .ok_or_else(|| { + exec_datafusion_err!("Timestamp is outside Chrono's calendar range") + }) + }; + let start = transition(start, standard_offset)?; + let end = transition(end, standard_offset + 3600)?; + Ok(if (start..end).contains(&seconds) { + standard_offset + 3600 + } else { + standard_offset + }) +} + +fn last_sunday(year: i32, month: u32, last_day: u32) -> Option { + let date = NaiveDate::from_ymd_opt(year, month, last_day)?; + date.with_day(last_day - date.weekday().num_days_from_sunday()) +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow::datatypes::TimestampMicrosecondType; + + fn shift(timezone: &str, timestamp: i64) -> i64 { + SparkTimeZone::parse(timezone) + .unwrap() + .adjust_to_local_time::(timestamp) + .unwrap() + } + + #[test] + fn spark_offset_spellings_preserve_sign_and_seconds() { + let offsets = [ + ("+0", 0), + ("-00", 0), + ("+1", 3600), + ("-01", -3600), + ("+0102", 3720), + ("-01:02", -3720), + ("+1:02", 3720), + ("+01:2", 3720), + ("+1:2", 3720), + ("+010203", 3723), + ("-01:02:03", -3723), + ("+1:02:03", 3723), + ("+00:00:01", 1), + ("-000001", -1), + ("+18:00", 64800), + ("-180000", -64800), + ]; + for prefix in ["", "UTC", "GMT", "UT"] { + for (offset, seconds) in offsets { + let timezone = format!("{prefix}{offset}"); + assert_eq!( + shift(&timezone, 123_456), + 123_456 + i64::from(seconds) * 1_000_000, + "{timezone}" + ); + } + } + for timezone in ["Z", "UTC", "GMT", "UT", "GMT0", "Etc/UTC"] { + assert_eq!(shift(timezone, -123_456), -123_456, "{timezone}"); + } + assert_eq!(shift("GMT+1", 0), 3_600_000_000); + assert_eq!(shift("Etc/GMT+1", 0), -3_600_000_000); + } + + #[test] + fn malformed_offsets_and_unknown_ids_are_rejected() { + for prefix in ["", "UTC", "GMT", "UT"] { + for offset in [ + "+", + "-", + "+19", + "-18:00:01", + "+18:01", + "+01:60", + "+01:99", + "+01:00:60", + "+1:2:03", + "+01:02:3", + "+123", + "+12345", + "+001:00", + "+01:", + "+:01", + "+01::00", + "+01:00:00:00", + "+01:00 ", + "+01", + "+01:٠0", + "++01", + "--01", + ] { + let timezone = format!("{prefix}{offset}"); + assert!(SparkTimeZone::parse(&timezone).is_err(), "{timezone}"); + } + } + for timezone in [ + "", + " ", + "utc", + "pst", + "PDT", + "ROC", + "America/Unknown", + "SystemV/PST9PDT", + " UTC", + ] { + assert!(SparkTimeZone::parse(timezone).is_err(), "{timezone}"); + } + } + + #[test] + fn java_short_ids_match_their_regions_in_winter_and_summer() { + let aliases = [ + ("ACT", "Australia/Darwin"), + ("AET", "Australia/Sydney"), + ("AGT", "America/Argentina/Buenos_Aires"), + ("ART", "Africa/Cairo"), + ("AST", "America/Anchorage"), + ("BET", "America/Sao_Paulo"), + ("BST", "Asia/Dhaka"), + ("CAT", "Africa/Harare"), + ("CNT", "America/St_Johns"), + ("CST", "America/Chicago"), + ("CTT", "Asia/Shanghai"), + ("EAT", "Africa/Addis_Ababa"), + ("ECT", "Europe/Paris"), + ("IET", "America/Indiana/Indianapolis"), + ("IST", "Asia/Kolkata"), + ("JST", "Asia/Tokyo"), + ("MIT", "Pacific/Apia"), + ("NET", "Asia/Yerevan"), + ("NST", "Pacific/Auckland"), + ("PLT", "Asia/Karachi"), + ("PNT", "America/Phoenix"), + ("PRT", "America/Puerto_Rico"), + ("PST", "America/Los_Angeles"), + ("SST", "Pacific/Guadalcanal"), + ("VST", "Asia/Ho_Chi_Minh"), + ("EST", "-05:00"), + ("HST", "-10:00"), + ("MST", "-07:00"), + ]; + for (alias, canonical) in aliases { + for timestamp in [1_704_067_200_123_456, 1_719_792_000_123_456] { + assert_eq!( + shift(alias, timestamp), + shift(canonical, timestamp), + "{alias}" + ); + } + } + } + + #[test] + fn system_v_zones_match_java_transition_boundaries() { + let zones = [ + ("SystemV/AST4ADT", -4), + ("SystemV/EST5EDT", -5), + ("SystemV/CST6CDT", -6), + ("SystemV/MST7MDT", -7), + ("SystemV/PST8PDT", -8), + ("SystemV/YST9YDT", -9), + ]; + for (timezone, hours) in zones { + let standard = hours * 3600; + for (local, offset_before, offset_after) in [ + ("1900-04-29T02:00:00Z", standard, standard + 3600), + ("1974-01-06T02:00:00Z", standard, standard + 3600), + ("1974-11-24T02:00:00Z", standard + 3600, standard), + ("1975-02-23T02:00:00Z", standard, standard + 3600), + ("2024-04-28T02:00:00Z", standard, standard + 3600), + ("2024-10-27T02:00:00Z", standard + 3600, standard), + ] { + let transition = DateTime::parse_from_rfc3339(local) + .unwrap() + .timestamp_micros() + - i64::from(offset_before) * 1_000_000; + assert_eq!( + shift(timezone, transition - 1), + transition - 1 + i64::from(offset_before) * 1_000_000, + "{timezone} before {local}" + ); + assert_eq!( + shift(timezone, transition), + transition + i64::from(offset_after) * 1_000_000, + "{timezone} at {local}" + ); + } + let before_transitions = DateTime::parse_from_rfc3339("1899-07-01T00:00:00Z") + .unwrap() + .timestamp_micros(); + assert_eq!( + shift(timezone, before_transitions), + before_transitions + i64::from(standard) * 1_000_000 + ); + } + for (timezone, hours) in [ + ("SystemV/AST4", -4), + ("SystemV/EST5", -5), + ("SystemV/CST6", -6), + ("SystemV/MST7", -7), + ("SystemV/PST8", -8), + ("SystemV/YST9", -9), + ("SystemV/HST10", -10), + ] { + assert_eq!( + shift(timezone, 1_719_792_000_000_000), + 1_719_792_000_000_000 + i64::from(hours) * 3_600_000_000 + ); + } + } +} diff --git a/datafusion/sqllogictest/test_files/spark/datetime/from_utc_timestamp.slt b/datafusion/sqllogictest/test_files/spark/datetime/from_utc_timestamp.slt index 8200fccfc1984..8b04956cb309e 100644 --- a/datafusion/sqllogictest/test_files/spark/datetime/from_utc_timestamp.slt +++ b/datafusion/sqllogictest/test_files/spark/datetime/from_utc_timestamp.slt @@ -154,3 +154,28 @@ query P SELECT from_utc_timestamp('2020-11-04T14:06:40'::timestamp, 'America/New_York'::string); ---- 2020-11-04T09:06:40 + +# Spark timezone spellings supplied as column values, including second-resolution offsets. +query P +SELECT from_utc_timestamp(column1, column2) +FROM VALUES +('2024-07-01T00:00:00.123456'::timestamp, 'PST'::string), +('2024-07-01T00:00:00.123456'::timestamp, 'GMT+1'::string), +('2024-07-01T00:00:00.123456'::timestamp, 'UTC+1:2'::string), +('2024-07-01T00:00:00.123456'::timestamp, 'UT-010203'::string), +('2024-04-15T00:00:00.123456'::timestamp, 'SystemV/PST8PDT'::string), +('2024-07-01T00:00:00.123456'::timestamp, 'SystemV/HST10'::string), +(NULL::timestamp, '+19:00'::string); +---- +2024-06-30T17:00:00.123456 +2024-07-01T01:00:00.123456 +2024-07-01T01:02:00.123456 +2024-06-30T22:57:57.123456 +2024-04-14T16:00:00.123456 +2024-06-30T14:00:00.123456 +NULL + +# Arrow accepts this offset, but Spark limits offsets to 18 hours. +query error invalid timezone '\+19:00' +SELECT from_utc_timestamp(column1, column2) +FROM VALUES ('2024-07-01'::timestamp, '+19:00'::string); diff --git a/typos.toml b/typos.toml index 196766f12fbc0..6a958616a9f17 100644 --- a/typos.toml +++ b/typos.toml @@ -17,6 +17,7 @@ Ois = "Ois" alo = "alo" # abbreviations, common words, etc. +IST = "IST" typ = "typ" datas = "datas" YOUY = "YOUY"