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
170 changes: 140 additions & 30 deletions datafusion/functions/src/datetime/date_bin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ use datafusion_expr::{
};
use datafusion_macros::user_doc;

use super::date_trunc::{civil_from_days, days_from_civil};
use chrono::{DateTime, Datelike, Duration, Months, TimeDelta, Utc};

#[user_doc(
Expand Down Expand Up @@ -463,24 +464,85 @@ fn date_bin_months_interval(stride_months: i64, source: i64, origin: i64) -> Res
}
}

// `date_bin_months_interval` with i128 nanoseconds, limited by the range of
// `DateTime<Utc>` instead of i64 nanoseconds.
// Compared in field order, which is chronological.
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
struct CivilTimestamp {
year: i64,
month: i64,
day: i64,
tod_nanos: i64,
}

// Month bins on the civil calendar. `DateTime<Utc>` cannot represent every
// second, millisecond, or microsecond timestamp (#25855).
fn date_bin_months_interval_wide(
stride_months: i64,
source: i128,
origin: i64,
) -> Option<i128> {
let nanos_per_sec = i128::from(NANOS_PER_SEC);
let secs = i64::try_from(source.div_euclid(nanos_per_sec)).ok()?;
let nsec = source.rem_euclid(nanos_per_sec) as u32;
let source_date = DateTime::from_timestamp(secs, nsec)?;
let origin_date = to_utc_date_time(origin).ok()?;

let bin_time = bin_months(stride_months, source_date, origin_date).ok()?;
Some(
i128::from(bin_time.timestamp()) * nanos_per_sec
+ i128::from(bin_time.timestamp_subsec_nanos()),
)
let source_ts = nanos_to_civil(source)?;
let origin_ts = nanos_to_civil(i128::from(origin))?;
let month_diff = month_index(source_ts)?.checked_sub(month_index(origin_ts)?)?;
let month_delta = compute_distance(month_diff, stride_months).ok()?;

// Step back one stride from the origin when the candidate is after the source.
let mut bin = shift_civil_months(origin_ts, month_delta)?;
if bin > source_ts {
let month_delta = month_delta.checked_sub(stride_months)?;
bin = shift_civil_months(origin_ts, month_delta)?;
}
civil_to_nanos(bin)
}

fn month_index(ts: CivilTimestamp) -> Option<i64> {
ts.year.checked_mul(12)?.checked_add(ts.month)
}

fn nanos_to_civil(nanos: i128) -> Option<CivilTimestamp> {
let nanos_per_day = i128::from(NANOSECONDS_IN_DAY);
let days = i64::try_from(nanos.div_euclid(nanos_per_day)).ok()?;
let tod_nanos = i64::try_from(nanos.rem_euclid(nanos_per_day)).ok()?;
let (year, month, day) = civil_from_days(days);
Some(CivilTimestamp {
year,
month,
day,
tod_nanos,
})
}

fn civil_to_nanos(ts: CivilTimestamp) -> Option<i128> {
let days = days_from_civil(ts.year, ts.month, ts.day);
i128::from(days)
.checked_mul(i128::from(NANOSECONDS_IN_DAY))?
.checked_add(i128::from(ts.tod_nanos))
}

// Keeps the origin's day and time. A missing day clamps to the month's end,
// matching `chrono`'s `checked_add_months`.
fn shift_civil_months(
origin: CivilTimestamp,
month_delta: i64,
) -> Option<CivilTimestamp> {
let total = (origin.month - 1).checked_add(month_delta)?;
let year = origin.year.checked_add(total.div_euclid(12))?;
let month = total.rem_euclid(12) + 1;
let last_day = days_in_month(year, month)?;
Some(CivilTimestamp {
year,
month,
day: origin.day.min(last_day),
tod_nanos: origin.tod_nanos,
})
}

fn days_in_month(year: i64, month: i64) -> Option<i64> {
let (next_year, next_month) = if month == 12 {
(year.checked_add(1)?, 1)
} else {
(year, month + 1)
};
days_from_civil(next_year, next_month, 1).checked_sub(days_from_civil(year, month, 1))
}

// return the start of the month bin that `source_date` falls into
Expand Down Expand Up @@ -942,7 +1004,10 @@ fn date_bin_impl(
mod tests {
use std::sync::Arc;

use crate::datetime::date_bin::{DateBinFunc, date_bin_nanos_interval};
use crate::datetime::date_bin::{
DateBinFunc, date_bin_months_interval, date_bin_months_interval_wide,
date_bin_nanos_interval,
};
use arrow::array::types::TimestampNanosecondType;
use arrow::array::{Array, IntervalDayTimeArray, TimestampNanosecondArray};
use arrow::compute::kernels::cast_utils::string_to_timestamp_nanos;
Expand Down Expand Up @@ -1507,14 +1572,22 @@ mod tests {
)),
];

let result = invoke_date_bin_with_args(args, 1, return_field);
assert!(result.is_ok());
if let ColumnarValue::Scalar(ScalarValue::TimestampMillisecond(val, _)) =
result.unwrap()
// The bin fits in milliseconds even though it is outside `DateTime<Utc>`.
let expected_ms = Some(-4_306_016_287_785_600_000);
let assert_millis = |args: Vec<ColumnarValue>| match invoke_date_bin_with_args(
args,
1,
return_field,
)
.unwrap()
{
assert!(val.is_none(), "Expected None for out of range operation");
}
let args = vec![
ColumnarValue::Scalar(ScalarValue::TimestampMillisecond(val, _)) => {
assert_eq!(val, expected_ms);
}
other => panic!("expected TimestampMillisecond, got {other:?}"),
};
assert_millis(args);
assert_millis(vec![
ColumnarValue::Scalar(ScalarValue::new_interval_mdn(1637426858, 0, 0)),
ColumnarValue::Scalar(ScalarValue::TimestampMillisecond(
Some(-1040292460),
Expand All @@ -1524,14 +1597,43 @@ mod tests {
Some(string_to_timestamp_nanos("1984-01-07 00:00:00").unwrap()),
None,
)),
];
]);
}

let result = invoke_date_bin_with_args(args, 1, return_field);
assert!(result.is_ok());
if let ColumnarValue::Scalar(ScalarValue::TimestampMillisecond(val, _)) =
result.unwrap()
{
assert!(val.is_none(), "Expected None for out of range operation");
#[test]
fn test_date_bin_months_wide_matches_chrono_path() {
let origins = [
0,
string_to_timestamp_nanos("2020-01-31T15:30:00").unwrap(),
string_to_timestamp_nanos("2020-02-29T23:59:59.500").unwrap(),
string_to_timestamp_nanos("1969-12-31T00:00:01").unwrap(),
string_to_timestamp_nanos("1970-03-01T00:00:00").unwrap(),
];
let sources = [
string_to_timestamp_nanos("2020-02-15T00:00:00").unwrap(),
string_to_timestamp_nanos("2020-02-29T00:00:00").unwrap(),
string_to_timestamp_nanos("2021-02-28T12:00:00").unwrap(),
string_to_timestamp_nanos("2021-03-01T00:00:00").unwrap(),
string_to_timestamp_nanos("2019-12-31T23:59:59").unwrap(),
string_to_timestamp_nanos("1960-01-31T00:00:00").unwrap(),
0,
-1,
];
for stride in [1, 2, 3, 12, 15, -1, -12] {
for &origin in &origins {
for &source in &sources {
let narrow =
date_bin_months_interval(stride, source, origin).unwrap();
let wide =
date_bin_months_interval_wide(stride, i128::from(source), origin)
.expect("wide month bin");
assert_eq!(
wide,
i128::from(narrow),
"stride={stride} source={source} origin={origin}"
);
}
}
}
}

Expand Down Expand Up @@ -1614,6 +1716,14 @@ mod tests {
(month(), secs(10_000_000_000), secs(month_start(2286, 11))),
// An exact bin at the minimum value.
(second(), secs(i64::MIN), secs(i64::MIN)),
// Regression for #25855: past `DateTime<Utc>`, the month bin still fits.
(month(), secs(10_000_000_000_000), secs(9_999_998_294_400)),
(
month(),
secs(-10_000_000_000_000),
secs(-10_000_001_059_200),
),
(month(), secs(i64::MAX), secs(9_223_372_036_854_460_800)),
];
for (stride, source, expected) in cases {
assert_date_bin(stride, source, epoch(), expected);
Expand Down Expand Up @@ -1642,10 +1752,10 @@ mod tests {
ScalarValue::new_interval_mdn(0, 0, 3),
ScalarValue::TimestampNanosecond(Some(i64::MIN), None),
),
// Month bins are limited by the range of `DateTime<Utc>`.
// The month bin of Timestamp(Second)::MIN does not fit (#25855).
(
ScalarValue::new_interval_mdn(1, 0, 0),
ScalarValue::TimestampSecond(Some(i64::MAX), None),
ScalarValue::TimestampSecond(Some(i64::MIN), None),
),
];
for (stride, source) in cases {
Expand Down
4 changes: 2 additions & 2 deletions datafusion/functions/src/datetime/date_trunc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -644,7 +644,7 @@ const DAYS_PER_ERA: i64 = 146_097;
/// This is a port of Howard Hinnant's `civil_from_days`, which documents the
/// derivation of the constants and the March-based year used below:
/// <https://howardhinnant.github.io/date_algorithms.html#civil_from_days>
fn civil_from_days(days: i64) -> (i64, i64, i64) {
pub(super) fn civil_from_days(days: i64) -> (i64, i64, i64) {
let z = days + DAYS_EPOCH_SHIFT;
let era = z.div_euclid(DAYS_PER_ERA);
let day_of_era = z.rem_euclid(DAYS_PER_ERA);
Expand All @@ -671,7 +671,7 @@ fn civil_from_days(days: i64) -> (i64, i64, i64) {
/// This is a port of Howard Hinnant's `days_from_civil`, which documents the
/// derivation of the constants and the March-based year used below:
/// <https://howardhinnant.github.io/date_algorithms.html#days_from_civil>
fn days_from_civil(year: i64, month: i64, day: i64) -> i64 {
pub(super) fn days_from_civil(year: i64, month: i64, day: i64) -> i64 {
let year = year - i64::from(month <= 2);
let era = year.div_euclid(400);
let year_of_era = year.rem_euclid(400);
Expand Down
53 changes: 35 additions & 18 deletions datafusion/sqllogictest/test_files/date_bin_errors.slt
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@

# Tests for DATE_BIN error handling with out-of-range values

# Test case from issue #20219 - should return NULL instead of panicking
query P
select date_bin(interval '1637426858 months', to_timestamp_millis(1040292460), timestamp '1984-01-07 00:00:00');
# Issue #20219: this stride must not panic. The bin fits in milliseconds (#25855).
query I
select arrow_cast(date_bin(interval '1637426858 months', to_timestamp_millis(1040292460), timestamp '1984-01-07 00:00:00'), 'Int64');
----
NULL
-4306016287785600000

# Issue #22528: negative sub-second source with month interval.
query P
Expand All @@ -36,12 +36,12 @@ from values (to_timestamp_millis(-1040292460)) t(c);
----
1969-12-07T00:00:00

# Array path should return NULL for per-row overflow.
query P
select date_bin(interval '1637426858 months', c, timestamp '1984-01-07 00:00:00')
# Array path matches the scalar path above.
query I
select arrow_cast(date_bin(interval '1637426858 months', c, timestamp '1984-01-07 00:00:00'), 'Int64')
from values (to_timestamp_millis(1040292460)) t(c);
----
NULL
-4306016287785600000

# Large stride causing overflow - should return NULL
query P
Expand All @@ -53,25 +53,42 @@ select date_bin(
----
NULL

# Another large stride test
query P
select date_bin(
# Same millisecond bin as the first query.
query I
select arrow_cast(date_bin(
interval '1637426858 months',
to_timestamp_millis(-1040292000),
timestamp '1984-01-07 00:00:00'
) as b;
), 'Int64');
----
NULL
-4306016287785600000

# Test with 1900-01-01 timestamp
query P
select date_bin(
# 1900-01-01 lands in that same bin.
query I
select arrow_cast(date_bin(
interval '1637426858 months',
to_timestamp_millis(-2208988800000),
timestamp '1984-01-07 00:00:00'
) as b;
), 'Int64');
----
NULL
-4306016287785600000

# Issue #25855: month stride past DateTime<Utc>, bin still fits.
query I
select arrow_cast(
date_bin(interval '1 month', arrow_cast(10000000000000, 'Timestamp(Second, None)')),
'Int64'
);
----
9999998294400

query I
select arrow_cast(
date_bin(interval '1 month', arrow_cast(-10000000000000, 'Timestamp(Second, None)')),
'Int64'
);
----
-10000001059200

# source - origin overflows i64 nanoseconds
query P
Expand Down