Skip to content
Merged
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
6 changes: 6 additions & 0 deletions vortex-array/proto/expr.proto
Original file line number Diff line number Diff line change
Expand Up @@ -129,3 +129,9 @@ message SelectOpts {
message CaseWhenOpts {
uint32 num_children = 1;
}

// Options for `vortex.replace_time_zone`. Ambiguity policy is an expression child.
message ReplaceTimeZoneOpts {
optional string time_zone = 1;
bool null_on_non_existent = 2;
}
26 changes: 26 additions & 0 deletions vortex-array/src/expr/exprs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ use crate::scalar_fn::fns::operators::CompareOperator;
use crate::scalar_fn::fns::operators::Operator;
use crate::scalar_fn::fns::pack::Pack;
use crate::scalar_fn::fns::pack::PackOptions;
use crate::scalar_fn::fns::replace_time_zone::ReplaceTimeZone;
use crate::scalar_fn::fns::replace_time_zone::ReplaceTimeZoneOptions;
use crate::scalar_fn::fns::select::FieldSelection;
use crate::scalar_fn::fns::select::Select;
use crate::scalar_fn::fns::variant_get::VariantGet;
Expand Down Expand Up @@ -1265,6 +1267,29 @@ pub fn bound_list_sum_opts(
.vortex_expect("list-sum expressions require a numeric list child")
}

/// Reinterpret timestamp wall times in `options.time_zone`, preserving the time unit.
///
/// `ambiguous` is a UTF-8 expression containing `raise`, `earliest`, `latest`, or `null`.
/// Removing the timezone preserves wall time in a timezone-naive timestamp.
pub fn replace_time_zone(
input: Expression,
ambiguous: Expression,
options: ReplaceTimeZoneOptions,
) -> Expression {
ReplaceTimeZone.new_expr(options, [input, ambiguous])
}

/// Create a bound timezone replacement expression.
pub fn bound_replace_time_zone(
input: BoundExpression,
ambiguous: BoundExpression,
options: ReplaceTimeZoneOptions,
) -> BoundExpression {
ReplaceTimeZone
.try_new_bound_expr(options, [input, ambiguous])
.vortex_expect("timezone replacement requires a timestamp and a UTF-8 ambiguity policy")
}

/// Constructors for expressions whose children have already been bound and type-checked.
///
/// These mirror the constructors in [`crate::expr`] and panic when the supplied children do not
Expand Down Expand Up @@ -1313,6 +1338,7 @@ pub mod bound {
pub use super::bound_or as or;
pub use super::bound_or_collect as or_collect;
pub use super::bound_pack as pack;
pub use super::bound_replace_time_zone as replace_time_zone;
pub use super::bound_root as root;
pub use super::bound_select as select;
pub use super::bound_select_exclude as select_exclude;
Expand Down
1 change: 1 addition & 0 deletions vortex-array/src/expr/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ pub use exprs::not_like;
pub use exprs::or;
pub use exprs::or_collect;
pub use exprs::pack;
pub use exprs::replace_time_zone;
pub use exprs::root;
pub use exprs::select;
pub use exprs::select_exclude;
Expand Down
37 changes: 34 additions & 3 deletions vortex-array/src/extension/datetime/timestamp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ use std::fmt;
use std::sync::Arc;

use jiff::Span;
use jiff::fmt::temporal::DateTimeParser;
use jiff::tz::TimeZone;
use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
Expand Down Expand Up @@ -75,6 +77,17 @@ impl fmt::Display for TimestampOptions {
}
}

/// Resolve an Arrow timestamp timezone, accepting IANA names and fixed offsets.
pub(crate) fn resolve_time_zone(name: Option<&str>) -> VortexResult<TimeZone> {
Ok(match name {
Some(name) if name.starts_with('+') || name.starts_with('-') => {
DateTimeParser::new().parse_time_zone(name)?
}
Some(name) => TimeZone::get(name)?,
None => TimeZone::UTC,
})
}

/// Unpacked value of a [`Timestamp`] extension scalar.
///
/// Each variant carries the raw storage value and an optional timezone.
Expand Down Expand Up @@ -102,7 +115,8 @@ impl fmt::Display for TimestampValue<'_> {
match tz {
None => write!(f, "{ts}"),
Some(tz) => {
let adjusted_ts = ts.in_tz(tz.as_ref()).vortex_expect("unknown timezone");
let zone = resolve_time_zone(Some(tz)).vortex_expect("unknown timezone");
let adjusted_ts = ts.to_zoned(zone);
write!(f, "{adjusted_ts}",)
}
}
Expand Down Expand Up @@ -218,12 +232,12 @@ impl ExtVTable for Timestamp {
};

// Validate the storage value is within the valid range for Timestamp.
let ts = jiff::Timestamp::UNIX_EPOCH
jiff::Timestamp::UNIX_EPOCH
.checked_add(span)
.map_err(|e| vortex_err!("Invalid timestamp scalar: {}", e))?;

if let Some(tz) = tz {
ts.in_tz(tz.as_ref())
resolve_time_zone(Some(tz))
.map_err(|e| vortex_err!("Invalid timezone for timestamp scalar: {}", e))?;
}

Expand All @@ -235,6 +249,7 @@ impl ExtVTable for Timestamp {
mod tests {
use std::sync::Arc;

use rstest::rstest;
use vortex_error::VortexResult;

use crate::dtype::DType;
Expand All @@ -253,6 +268,22 @@ mod tests {
Ok(())
}

#[cfg_attr(miri, ignore)]
#[rstest]
#[case("+01:00", "1970-01-01T01:00:00+01:00[+01:00]")]
#[case("-05:30", "1969-12-31T18:30:00-05:30[-05:30]")]
fn display_fixed_offset_timestamp(
#[case] timezone: &str,
#[case] expected: &str,
) -> VortexResult<()> {
let dtype = DType::Extension(
Timestamp::new_with_tz(TimeUnit::Seconds, Some(timezone.into()), Nullable).erased(),
);
let scalar = Scalar::try_new(dtype, Some(ScalarValue::Primitive(PValue::I64(0))))?;
assert_eq!(format!("{}", scalar.as_extension()), expected);
Ok(())
}

#[cfg_attr(miri, ignore)]
#[test]
fn reject_timestamp_with_invalid_timezone() {
Expand Down
8 changes: 8 additions & 0 deletions vortex-array/src/proto/generated/vortex.expr.rs
Original file line number Diff line number Diff line change
Expand Up @@ -217,3 +217,11 @@ pub struct CaseWhenOpts {
#[prost(uint32, tag = "1")]
pub num_children: u32,
}
/// Options for `vortex.replace_time_zone`. Ambiguity policy is an expression child.
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct ReplaceTimeZoneOpts {
#[prost(string, optional, tag = "1")]
pub time_zone: ::core::option::Option<::prost::alloc::string::String>,
#[prost(bool, tag = "2")]
pub null_on_non_existent: bool,
}
1 change: 1 addition & 0 deletions vortex-array/src/scalar_fn/fns/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ pub mod merge;
pub mod not;
pub mod operators;
pub mod pack;
pub mod replace_time_zone;
pub mod select;
pub mod stat;
pub mod variant_get;
Expand Down
Loading
Loading