From d9d943c93e7344c350ce0ae44271a125c60b85ac Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 07:14:38 -0600 Subject: [PATCH 1/6] fix: match Iceberg's rounding for pre-1970 timestamps in native years/months/days/hours Iceberg's DateTimeUtil.convertMicros does not quite floor. For a negative timestamp it builds an instant from floorDiv(micros, 1_000_000) seconds and floorMod(micros + 1, 1_000_000) microseconds, and when the microsecond of second is 999999 the added microsecond wraps without carrying into the seconds. A timestamp exactly 999999 microseconds into a unit therefore gets the unit before: Iceberg puts 1969-01-01 00:00:00.999999 in year -2, month -13, day 1968-12-31 and hour -8761. Its Spark functions return those values and its partition transforms write them. The native kernels floored and returned -1, -12, 1969-01-01 and -8760. The timestamp kernels now take one unit off in that case. Dates and non-negative timestamps are unchanged. The native Iceberg writer computes year and month partition values with these kernels. It computed day and hour with iceberg-rust's transforms, which floor, and whose day also moves some timestamps from the last second before a pre-1970 midnight into the next day. It now computes day and hour of a timestamp with the kernels too, so its partition values match iceberg-java's and the sort in front of a clustered write. Closes #6426. --- .github/workflows/pr_build_linux.yml | 1 + .github/workflows/pr_build_macos.yml | 1 + .../contributor-guide/iceberg-writes.md | 15 +- .../operators/iceberg_partition_value.rs | 138 ++++++++++++++---- .../src/execution/operators/iceberg_write.rs | 27 ++-- .../spark-expr/src/iceberg_funcs/temporal.rs | 130 ++++++++++++++++- ...IcebergSystemFunctionExtensionsSuite.scala | 81 ++++++++++ .../CometIcebergSystemFunctionSuite.scala | 101 ++++++++++++- 8 files changed, 441 insertions(+), 53 deletions(-) create mode 100644 spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index ee8a947aa47..275613fa296 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -518,6 +518,7 @@ jobs: org.apache.comet.CometIcebergWriteActionSuite org.apache.comet.CometIcebergWriteDetectionSuite org.apache.comet.CometIcebergSystemFunctionSuite + org.apache.comet.CometIcebergSystemFunctionExtensionsSuite org.apache.comet.CometIcebergResidualPushdownSuite org.apache.comet.iceberg.IcebergReflectionSuite org.apache.comet.serde.operator.IcebergWriteProtoTranslationSuite diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 371e4f6531c..e695635ed08 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -175,6 +175,7 @@ jobs: org.apache.comet.CometIcebergWriteActionSuite org.apache.comet.CometIcebergWriteDetectionSuite org.apache.comet.CometIcebergSystemFunctionSuite + org.apache.comet.CometIcebergSystemFunctionExtensionsSuite org.apache.comet.CometIcebergResidualPushdownSuite org.apache.comet.iceberg.IcebergReflectionSuite org.apache.comet.serde.operator.IcebergWriteProtoTranslationSuite diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index b30720f9e92..d3a51ed4702 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -215,11 +215,16 @@ Points where Comet adapts iceberg-rust to match iceberg-java: `LocationGenerator::generate_location` cannot return an error. - **Partition values.** `PartitionValueCalculator` (`iceberg_partition_value.rs`) replaces iceberg-rust's calculator of the same name, and `PartitionSplitter` replaces its - `RecordBatchPartitionSplitter`. `year` and `month` go through Comet's `iceberg_years` / - `iceberg_months` kernels, the ones the sort in front of a clustered write runs. iceberg-rust - computes them with Arrow's `date_part`, which returns NULL past `chrono`'s calendar - ([#6145](https://github.com/apache/datafusion-comet/issues/6145)). Every other transform stays on - iceberg-rust. + `RecordBatchPartitionSplitter`. `year`, `month`, `day` and `hour` of a date or timestamp go + through Comet's `iceberg_years` / `iceberg_months` / `iceberg_days` / `iceberg_hours` kernels, the + ones the sort in front of a clustered write runs (`day` of a date is the date itself and stays on + iceberg-rust). iceberg-rust computes `year` and `month` with Arrow's `date_part`, which returns + NULL past `chrono`'s calendar ([#6145](https://github.com/apache/datafusion-comet/issues/6145)). + All four of its transforms floor a pre-1970 timestamp that lies exactly 999999 microseconds into a + unit, which iceberg-java puts in the unit before + ([#6426](https://github.com/apache/datafusion-comet/issues/6426)), and its `day` puts some + timestamps from the last second before a pre-1970 midnight in the next day. Every other transform + stays on iceberg-rust. - **File names.** `file_name_prefix` embeds the partition id, the task attempt id and the operation id, so a retried or speculative attempt never reuses another attempt's file names. - **Row pacing.** iceberg-java's rolling writer checks the target file size every 1000 rows of the diff --git a/native/core/src/execution/operators/iceberg_partition_value.rs b/native/core/src/execution/operators/iceberg_partition_value.rs index a4767cdb5cb..31b724a3e89 100644 --- a/native/core/src/execution/operators/iceberg_partition_value.rs +++ b/native/core/src/execution/operators/iceberg_partition_value.rs @@ -19,8 +19,8 @@ //! //! [`PartitionValueCalculator`] stands in for iceberg-rust's calculator of the same name. It //! projects each partition field's source column and applies the field's transform the same way, -//! but computes `year` and `month` with Comet's own kernels, so that a date or timestamp past -//! `chrono`'s calendar gets iceberg-java's partition value instead of a NULL. +//! but computes the time transforms of dates and timestamps with Comet's own kernels, so that they +//! get iceberg-java's partition values where iceberg-rust's differ. use std::sync::Arc; @@ -35,20 +35,31 @@ use iceberg::{Error, ErrorKind, Result}; /// Computes a batch's partition values: one row of the partition struct per input row. /// -/// Matches iceberg-rust's `PartitionValueCalculator` except for `year` and `month` over a `date`, -/// `timestamp`, or `timestamptz` source. iceberg-rust splits the calendar with Arrow's `date_part`, -/// which returns NULL for anything `chrono` cannot represent -- past year 262142 -- whereas -/// iceberg-java's `DateTimeUtil` goes through `LocalDate` and covers every Spark date (to year -/// 5881580) and timestamp (to year 294247). The NULL did not fail the write: the data file was -/// committed claiming a NULL partition for rows whose source value is not NULL -/// (apache/datafusion-comet#6145). Those two transforms go through Comet's `iceberg_years` / -/// `iceberg_months` kernels instead, the ones the sort in front of a clustered write runs, which -/// are pinned against iceberg-java over the whole domain. Wherever `chrono` can represent the date -/// the two implementations agree, so every value iceberg-rust could compute is unchanged. +/// Matches iceberg-rust's `PartitionValueCalculator` except for the time transforms of a `date`, +/// `timestamp`, or `timestamptz` source, which go through Comet's `iceberg_years` / +/// `iceberg_months` / `iceberg_days` / `iceberg_hours` kernels instead: the ones the sort in front +/// of a clustered write runs, pinned against iceberg-java's `DateTimeUtil` over the whole domain. +/// iceberg-rust's transforms differ from iceberg-java's in three ways: /// -/// `day` and `hour` stay on iceberg-rust: they are floor divisions of the epoch value and never -/// consult the calendar. So do the nanosecond timestamp types, whose `i64` range (years 1677 to -/// 2262) lies inside `chrono`'s and which Comet's kernels do not accept. +/// - `year` and `month` split the calendar with Arrow's `date_part`, which returns NULL for +/// anything `chrono` cannot represent -- past year 262142 -- whereas iceberg-java goes through +/// `LocalDate` and covers every Spark date (to year 5881580) and timestamp (to year 294247). The +/// NULL did not fail the write: the data file was committed claiming a NULL partition for rows +/// whose source value is not NULL (apache/datafusion-comet#6145). +/// - All four floor a pre-epoch timestamp that lies exactly 999999 microseconds into a unit, which +/// iceberg-java puts in the unit before, so `1969-01-01T00:00:00.999999` belongs in the 1968 +/// partitions (apache/datafusion-comet#6426). +/// - `day` moves a timestamp from the last second of a day before 1969-12-31 into the next day, +/// unless its microsecond of second is 0 or 999999: it takes the whole seconds with a truncating +/// division and the microseconds with a flooring one. +/// +/// A partition value that differs from the sort key can also fail a clustered write, which rejects +/// a row whose partition it has already closed. Everywhere else the two implementations agree, so +/// no other value changes. +/// +/// `day` of a `date` is the date itself and stays on iceberg-rust. So do the nanosecond timestamp +/// types, whose `i64` range (years 1677 to 2262) lies inside `chrono`'s and which Comet's kernels +/// do not accept. pub(crate) struct PartitionValueCalculator { projector: RecordBatchProjector, transforms: Vec, @@ -125,12 +136,14 @@ enum FieldTransform { impl FieldTransform { fn try_new(transform: Transform, source_type: Option<&Type>) -> Result { - let calendar_source = matches!( + let timestamp_source = matches!( source_type, Some(Type::Primitive( - PrimitiveType::Date | PrimitiveType::Timestamp | PrimitiveType::Timestamptz + PrimitiveType::Timestamp | PrimitiveType::Timestamptz )) ); + let calendar_source = + timestamp_source || matches!(source_type, Some(Type::Primitive(PrimitiveType::Date))); Ok(match transform { Transform::Year if calendar_source => { Self::Comet(SparkIcebergTemporalTransform::years()) @@ -138,6 +151,12 @@ impl FieldTransform { Transform::Month if calendar_source => { Self::Comet(SparkIcebergTemporalTransform::months()) } + Transform::Day if timestamp_source => { + Self::Comet(SparkIcebergTemporalTransform::days()) + } + Transform::Hour if timestamp_source => { + Self::Comet(SparkIcebergTemporalTransform::hours()) + } _ => Self::IcebergRust(create_transform_function(&transform)?), }) } @@ -328,10 +347,11 @@ mod tests { } } - // The transforms left on iceberg-rust already agree with iceberg-java this far out (same JVM - // run as above). `day` of a date is the date itself, so only timestamps are interesting. + // iceberg-rust's `day` and `hour` already agreed with iceberg-java this far out, and the Comet + // kernels that now compute them for timestamps still do (same JVM run as above). `day` of a + // date is the date itself, so only timestamps are interesting. #[test] - fn days_and_hours_past_chronos_calendar_already_match_iceberg_java() { + fn days_and_hours_past_chronos_calendar_match_iceberg_java() { let instants = micros_utc(&INSTANTS_PAST_CHRONO); let (comet, _, batch) = calculators(&[ (Transform::Day, Arc::clone(&instants)), @@ -362,9 +382,73 @@ mod tests { ); } - /// Every value iceberg-rust could already compute comes out unchanged, up to both ends of - /// `chrono`'s calendar, so what the native writer puts in a table now is consistent with what - /// it put there before. + // Expectations from iceberg-java 1.11's `DateTimeUtil` on a JDK 17 JVM; 1.5.2, 1.8.1, and + // 1.10.0 agree. + #[test] + fn pre_epoch_timestamps_partition_like_iceberg_java() { + let values = [ + // 1969-01-01T00:00:00.999999, which iceberg-java places by the second before it, + // 1968-12-31T23:59:59 (apache/datafusion-comet#6426). + Some(-31_535_999_000_001), + // 1969-12-31T23:00:00.999999, where that moves only the hour. + Some(-3_599_000_001), + // 1969-12-30T23:59:59.5 and 1969-12-30T23:59:59.999998. + Some(-86_400_500_000), + Some(-86_400_000_002), + None, + ]; + let day = |days: [i32; 4]| -> ArrayRef { + Arc::new(Date32Array::from_iter( + days.into_iter().map(Some).chain([None]), + )) + }; + let int = |values: [i32; 4]| -> ArrayRef { Arc::new(ints(&values)) }; + // (transform, iceberg-java's partition values, iceberg-rust's) + let cases = [ + ( + Transform::Year, + int([-2, -1, -1, -1]), + int([-1, -1, -1, -1]), + ), + ( + Transform::Month, + int([-13, -1, -1, -1]), + int([-12, -1, -1, -1]), + ), + ( + Transform::Day, + day([-366, -1, -2, -2]), + day([-365, -1, -1, -1]), + ), + ( + Transform::Hour, + int([-8_761, -2, -25, -25]), + int([-8_760, -1, -25, -25]), + ), + ]; + for source in [micros(&values), micros_utc(&values)] { + let fields: Vec<_> = cases + .iter() + .map(|(transform, _, _)| (*transform, Arc::clone(&source))) + .collect(); + let (comet, iceberg_rust, batch) = calculators(&fields); + let comet = columns(comet.calculate(&batch).unwrap()); + let iceberg_rust = columns(iceberg_rust.calculate(&batch).unwrap()); + for (i, (transform, java, rust)) in cases.iter().enumerate() { + let label = format!("{transform} of {}", source.data_type()); + assert_eq!(&comet[i], java, "{label}"); + // The reason Comet computes these: iceberg-rust floors the first two rows, and its + // `day` moves the last two into 1969-12-31. If this starts failing, iceberg-rust's + // transforms have changed and delegating needs another look. + assert_eq!(&iceberg_rust[i], rust, "iceberg-rust's {label}"); + } + } + } + + /// Away from the values where iceberg-rust parts from iceberg-java (the NULLs past `chrono`'s + /// calendar and the pre-epoch timestamps above), every value comes out as iceberg-rust + /// computes it, up to both ends of `chrono`'s calendar, so what the native writer puts in a + /// table now is consistent with what it put there before. #[test] fn agrees_with_iceberg_rust_wherever_chrono_can_represent_the_date() { let days = dates(&[ @@ -421,13 +505,17 @@ mod tests { ]; let (comet, iceberg_rust, batch) = calculators(&fields); - // The first four really do run through Comet's kernels, and nothing else does. + // The time transforms of dates and timestamps really do run through Comet's kernels, except + // `day` of a date, and nothing else does. let on_comet = comet .transforms .iter() .map(|transform| matches!(transform, FieldTransform::Comet(_))) .collect::>(); - assert_eq!(on_comet, [[true; 4].as_slice(), &[false; 9]].concat()); + assert_eq!( + on_comet, + [true, true, true, true, false, false, true, true, false, false, false, false, false] + ); for ((comet, iceberg_rust), (transform, source)) in columns(comet.calculate(&batch).unwrap()) diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index e4315c0f4de..baef664096f 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -3314,11 +3314,11 @@ mod tests { /// A partitioned write runs both: the sort in front of [`IcebergWriteExec`] is keyed on the /// `datafusion-comet-spark-expr` kernels (Iceberg plans the sort as `bucket(...)`, `days(...)`, /// ... system-function calls), while [`ClusteredWriter`] groups the sorted rows by the partition -/// values that [`PartitionValueCalculator`] computes -- with the same kernels for `years` and -/// `months`, and with iceberg-rust's transforms for everything else. The writer requires the two -/// to agree: when they do not it fails at runtime with "The input is not sorted! Cannot write to -/// partition that was previously closed". These tests make an iceberg-rust bump that changes a -/// transform break here first. +/// values that [`PartitionValueCalculator`] computes -- with the same kernels for the time +/// transforms of dates and timestamps, and with iceberg-rust's transforms for everything else. The +/// writer requires the two to agree: when they do not it fails at runtime with "The input is not +/// sorted! Cannot write to partition that was previously closed". These tests make an iceberg-rust +/// bump that changes a transform break here first. #[cfg(test)] mod iceberg_rust_transform_parity { use arrow::array::{ @@ -3569,7 +3569,10 @@ mod iceberg_rust_transform_parity { } } - /// `days` and `hours` are plain floor division on both sides, so the whole domain agrees. + /// `days` and `hours` agree with iceberg-rust except on the pre-epoch timestamps where + /// iceberg-rust parts from iceberg-java (see `PartitionValueCalculator`), which is why the + /// writer computes both for timestamps with the Comet kernels; agreeing here means the switch + /// changed no other value. `day` of a date still goes through iceberg-rust. #[test] fn days_and_hours_agree_with_iceberg_rust() { let micros = vec![ @@ -3603,12 +3606,12 @@ mod iceberg_rust_transform_parity { assert_agree("days(date)", Transform::Day, &days_udf, &dates); } - /// `years` and `months` agree over the dates iceberg-rust can represent -- it splits the - /// calendar with `chrono`, so anything past year 262142 comes back NULL there while Comet and - /// the JVM keep going (apache/iceberg-rust#3142; see the kernel's own unit tests for those). - /// That is why the writer computes these two with the Comet kernels itself - /// (apache/datafusion-comet#6145); agreeing here means the switch changed no value iceberg-rust - /// could compute. + /// `years` and `months` agree over the dates iceberg-rust can represent, apart from the + /// pre-epoch timestamps where it parts from iceberg-java (see `PartitionValueCalculator`). It + /// splits the calendar with `chrono`, so anything past year 262142 comes back NULL there while + /// Comet and the JVM keep going (apache/iceberg-rust#3142; see the kernel's own unit tests for + /// those). That is why the writer computes these two with the Comet kernels itself + /// (apache/datafusion-comet#6145); agreeing here means the switch changed no other value. #[test] fn years_and_months_agree_with_iceberg_rust_within_its_range() { let years_udf = SparkIcebergTemporalTransform::years(); diff --git a/native/spark-expr/src/iceberg_funcs/temporal.rs b/native/spark-expr/src/iceberg_funcs/temporal.rs index 96a2ee82d0b..717e72f48c6 100644 --- a/native/spark-expr/src/iceberg_funcs/temporal.rs +++ b/native/spark-expr/src/iceberg_funcs/temporal.rs @@ -19,8 +19,10 @@ //! //! Iceberg's `DateTimeUtil` evaluates all four in UTC regardless of the Spark session timezone //! (`TimestampType` and `TimestampNTZType` are handled identically), and all four floor: a value -//! before the epoch maps to a negative period. `years` and `months` are calendar-aware, `days` and -//! `hours` are plain floor division of the epoch value. `days` returns a date (Iceberg's +//! before the epoch maps to a negative period. The one exception is a pre-epoch timestamp whose +//! microsecond of second is 999999 right after a unit boundary, which Iceberg puts in the unit +//! before (see `iceberg_div_floor`). `years` and `months` are calendar-aware, `days` and `hours` +//! are floor division of the epoch value. `days` returns a date (Iceberg's //! `DaysFunction.resultType()` is `DateType`), the other three return an int. //! //! The kernels read the raw epoch values instead of Arrow's timezone-aware `date_part`. That is a @@ -41,9 +43,10 @@ use datafusion::common::{utils::take_function_args, Result}; use datafusion::logical_expr::{ ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl, Signature, Volatility, }; -use num::integer::div_floor; +use num::integer::{div_floor, mod_floor}; use std::sync::Arc; +const MICROS_PER_SECOND: i64 = 1_000_000; const MICROS_PER_HOUR: i64 = 3_600_000_000; const MICROS_PER_DAY: i64 = 86_400_000_000; const UNIX_EPOCH_YEAR: i32 = 1970; @@ -74,18 +77,41 @@ impl TemporalUnit { } } -/// `DateTimeUtil.microsToDays`: floor division, so `-1` micros is day `-1`. The quotient of the -/// widest `i64` micros is about 1.07e8, so the narrowing is always exact here. +/// `micros` divided by `unit_micros`, a whole number of seconds, the way Iceberg's +/// `DateTimeUtil.convertMicros` divides it. +/// +/// That is a floor, except for a negative value that lies exactly 999999 microseconds into a unit, +/// which gets the unit before. For a negative value Java builds an instant from +/// `floorDiv(micros, 1_000_000)` seconds and `floorMod(micros + 1, 1_000_000)` microseconds, and +/// subtracts one from the number of whole units between the epoch and that instant. The added +/// microsecond is what makes the count a floor, but at 999999 it wraps to 0 without carrying into +/// the seconds, so the instant is a second early, and when that second starts a unit the count +/// comes out one short. Iceberg puts `1969-01-01T00:00:00.999999` in year -2, month -13, day +/// 1968-12-31, and hour -8761 rather than -1, -12, 1969-01-01, and -8760. Its Spark functions +/// return those values and its partition transforms write them, from 1.5.2 through 1.11.0. +#[inline] +fn iceberg_div_floor(micros: i64, unit_micros: i64) -> i64 { + let units = div_floor(micros, unit_micros); + if micros < 0 && mod_floor(micros, unit_micros) == MICROS_PER_SECOND - 1 { + units - 1 + } else { + units + } +} + +/// `DateTimeUtil.microsToDays`. The quotient of the widest `i64` micros is about 1.07e8, so the +/// narrowing is always exact here. A month or year boundary is also a day boundary, so the +/// calendar split of this day gives `microsToMonths` and `microsToYears` too. #[inline] fn micros_to_days(micros: i64) -> i32 { - div_floor(micros, MICROS_PER_DAY) as i32 + iceberg_div_floor(micros, MICROS_PER_DAY) as i32 } /// `DateTimeUtil.microsToHours`. Java narrows the hour count with a plain `(int)` cast, which /// wraps beyond about 7.7e18 micros; `as i32` truncates the same way. #[inline] fn micros_to_hours(micros: i64) -> i32 { - div_floor(micros, MICROS_PER_HOUR) as i32 + iceberg_div_floor(micros, MICROS_PER_HOUR) as i32 } /// `DateTimeUtil.daysToYears`: whole calendar years between the epoch and the day, floored. @@ -231,6 +257,32 @@ mod tests { .unwrap() } + /// `[years, months, days, hours]` of each of `micros`, as a timestamp column tagged `tz`. + fn all_four(micros: &[i64], tz: Option<&str>) -> Vec<[i32; 4]> { + let mut array = TimestampMicrosecondArray::from(micros.to_vec()); + if let Some(tz) = tz { + array = array.with_timezone(tz); + } + let input: ArrayRef = Arc::new(array); + let [years, months, days, hours] = [ + TemporalUnit::Years, + TemporalUnit::Months, + TemporalUnit::Days, + TemporalUnit::Hours, + ] + .map(|unit| transform(unit, Arc::clone(&input))); + (0..micros.len()) + .map(|i| { + [ + years.as_primitive::().value(i), + months.as_primitive::().value(i), + days.as_primitive::().value(i), + hours.as_primitive::().value(i), + ] + }) + .collect() + } + // Boundaries around the epoch, as (epoch days, years, months). The values past the epoch // block are outside `chrono::NaiveDate`'s range but well inside `LocalDate`'s; they come from // running Iceberg's `DateTimeUtil.convertDays` on a JDK 17 JVM. @@ -307,6 +359,15 @@ mod tests { // `(int)` narrowing wraps it exactly as `as i32` does. (i64::MAX, 292_277, 3_507_324, 106_751_991, -1_732_919_508), (i64::MIN, -292_278, -3_507_325, -106_751_992, 1_732_919_507), + // The lowest `i64` that ends in 999999, where moving the value a second earlier would + // overflow. + ( + i64::MIN + 775_807, + -292_278, + -3_507_325, + -106_751_992, + 1_732_919_507, + ), ( 8_000_000_000_000_000_000, 253_509, @@ -373,6 +434,61 @@ mod tests { } } + /// Unit boundaries on both sides of the epoch, as (epoch second, and Iceberg's `[years, + /// months, days, hours]` for the last microsecond before the boundary and for the boundary + /// itself). The values come from `DateTimeUtil` on a JDK 17 JVM; Iceberg 1.5.2, 1.8.1, 1.10.0, + /// and 1.11.0 agree on them. + const BOUNDARIES: &[(i64, [i32; 4], [i32; 4])] = &[ + // 1969-12-31T23:59:30, a whole second that starts no unit. + (-30, [-1, -1, -1, -1], [-1, -1, -1, -1]), + // 1969-12-31T23:00:00, 1969-12-31, 1969-12-01, 1969-01-01, and 1900-01-01. + (-3_600, [-1, -1, -1, -2], [-1, -1, -1, -1]), + (-86_400, [-1, -1, -2, -25], [-1, -1, -1, -24]), + (-2_678_400, [-1, -2, -32, -745], [-1, -1, -31, -744]), + ( + -31_536_000, + [-2, -13, -366, -8_761], + [-1, -12, -365, -8_760], + ), + ( + -2_208_988_800, + [-71, -841, -25_568, -613_609], + [-70, -840, -25_567, -613_608], + ), + // The epoch, then 1970-01-01T01:00:00, 1970-01-02, 1970-02-01, and 1971-01-01. + (0, [-1, -1, -1, -1], [0, 0, 0, 0]), + (3_600, [0, 0, 0, 0], [0, 0, 0, 1]), + (86_400, [0, 0, 0, 23], [0, 0, 1, 24]), + (2_678_400, [0, 0, 30, 743], [0, 1, 31, 744]), + (31_536_000, [0, 11, 364, 8_759], [1, 12, 365, 8_760]), + ]; + + /// A pre-epoch timestamp whose microsecond of second is 999999 takes the unit of the second + /// before it, so one that follows a boundary lands in the unit before the boundary + /// (apache/datafusion-comet#6426). Its neighbours, and every timestamp from the epoch on, + /// floor. + #[test] + fn timestamps_ending_in_999999_match_iceberg_date_time_util() { + for &(second, before, at) in BOUNDARIES { + let boundary = second * MICROS_PER_SECOND; + let cases = [ + (boundary - 1, before), + (boundary, at), + (boundary + 999_998, at), + (boundary + 999_999, if boundary < 0 { before } else { at }), + (boundary + 1_000_000, at), + (boundary + 1_999_999, at), + ]; + let micros = cases.map(|(micros, _)| micros); + // The two tags Comet produces: untagged `TimestampNTZType` and `UTC` `TimestampType`. + for tz in [None, Some("UTC")] { + for ((micros, expected), actual) in cases.iter().zip(all_four(µs, tz)) { + assert_eq!(actual, *expected, "{micros} micros, {tz:?}"); + } + } + } + } + /// The pinned cases above check the endpoints of the domain; this checks the calendar split /// itself everywhere `chrono` can still represent the date, so the hand-rolled arithmetic /// cannot drift in between. diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala new file mode 100644 index 00000000000..09b563e929b --- /dev/null +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala @@ -0,0 +1,81 @@ +/* + * 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. + */ + +package org.apache.comet + +import org.apache.spark.SparkConf +import org.apache.spark.sql.CometTestBase +import org.apache.spark.sql.catalyst.expressions.ApplyFunctionExpression +import org.apache.spark.sql.catalyst.plans.logical.Filter +import org.apache.spark.sql.internal.SQLConf + +/** + * Iceberg's system functions once Iceberg's SQL extensions have rewritten them. + * + * The extensions' `ReplaceStaticInvoke` rule turns a system-function call that a filter compares + * with a constant from a `StaticInvoke` into an `ApplyFunctionExpression`, so that Iceberg can + * push the comparison into its scan. Over any other source the filter stays, and Comet evaluates + * it with the same native kernels as the `StaticInvoke` that CometIcebergSystemFunctionSuite + * covers. `spark.sql.extensions` is static, so this path needs a suite of its own. + */ +class CometIcebergSystemFunctionExtensionsSuite extends CometTestBase with CometIcebergTestBase { + + override protected def sparkConf: SparkConf = + super.sparkConf.set( + "spark.sql.extensions", + "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") + + test("rewritten temporal filters match Iceberg on pre-1970 timestamps ending in .999999") { + assume(icebergAvailable, "Iceberg not available in classpath") + withTempIcebergDir { warehouseDir => + withSQLConf( + "spark.sql.catalog.ice" -> "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.ice.type" -> "hadoop", + "spark.sql.catalog.ice.warehouse" -> warehouseDir.getAbsolutePath, + SQLConf.SESSION_LOCAL_TIMEZONE.key -> "UTC", + SQLConf.PARQUET_OUTPUT_TIMESTAMP_TYPE.key -> "TIMESTAMP_MICROS") { + withTable("pre_epoch") { + // A parquet table, so that no scan absorbs the filter and Comet evaluates it. Iceberg + // places a pre-1970 timestamp ending in .999999 by the second before it: the first row + // is in hour -8761, day 1968-12-31, month -13, and year -2, and the second in hour -2. + sql("CREATE TABLE pre_epoch (id INT, ts TIMESTAMP) USING parquet") + sql("""INSERT INTO pre_epoch VALUES + (1, TIMESTAMP '1969-01-01 00:00:00.999999'), + (2, TIMESTAMP '1969-12-31 23:00:00.999999'), + (3, TIMESTAMP '1969-12-31 22:30:00'), + (4, TIMESTAMP '1968-12-31 12:00:00')""") + Seq( + "ice.system.hours(ts) = -2", + "ice.system.days(ts) = DATE '1968-12-31'", + "ice.system.months(ts) = -13", + "ice.system.years(ts) = -2").foreach { predicate => + val df = sql(s"SELECT id FROM pre_epoch WHERE $predicate") + val plan = df.queryExecution.optimizedPlan + assert( + plan + .collect { case filter: Filter => filter.condition } + .exists(_.find(_.isInstanceOf[ApplyFunctionExpression]).isDefined), + s"expected Iceberg's extensions to rewrite $predicate:\n$plan") + checkSparkAnswerAndOperator(df) + } + } + } + } + } +} diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala index 118b2371d53..bf25e23e188 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala @@ -174,6 +174,24 @@ class CometIcebergSystemFunctionSuite } } + test("years, months, days, and hours match Iceberg on pre-1970 timestamps ending in .999999") { + // The random corpus almost never lands on such a value, and its boundary rows sit a + // microsecond before a boundary rather than after it. + withIcebergCatalog { + withPreEpochTable { + val transformed = for { + column <- Seq("ts", "ntz") + function <- Seq("years", "months", "days", "hours") + } yield s"$catalog.system.$function($column)" + checkSparkAnswerAndOperator(s"SELECT id, ${transformed.mkString(", ")} FROM pre_epoch") + checkSparkAnswerAndOperator( + s"SELECT id FROM pre_epoch WHERE $catalog.system.days(ts) = DATE '1968-12-31'") + checkSparkAnswerAndOperator( + s"SELECT id FROM pre_epoch WHERE $catalog.system.hours(ts) = -2") + } + } + } + test("system functions in filters stay native") { withSourceTable { checkSparkAnswerAndOperator( @@ -263,6 +281,44 @@ class CometIcebergSystemFunctionSuite } } + test("native partitioned write puts pre-1970 timestamps ending in .999999 where Iceberg does") { + // The sort in front of the write runs the native kernels and the native writer computes each + // row's partition value, so both have to follow Iceberg for the table to hold the partitions + // iceberg-java would have written. A spec takes one time transform per source column, hence a + // column per transform. + withIcebergCatalog { + withPreEpochTable { + val table = s"$catalog.db.pre_epoch_partitions" + sql(s""" + CREATE TABLE $table (id INT, y TIMESTAMP, m TIMESTAMP, d TIMESTAMP, h TIMESTAMP) + USING iceberg + PARTITIONED BY (years(y), months(m), days(d), hours(h))""") + try { + val plans = capturePlans(spark) { + sql(s"INSERT INTO $table SELECT id, ts, ts, ts, ts FROM pre_epoch") + } + assert( + plans.exists(plan => + collectWithSubqueries(plan) { case w: CometIcebergWriteExec => w }.nonEmpty), + s"expected a native Iceberg write in the captured plans:\n${plans.mkString("\n--\n")}") + + withSQLConf(CometConf.COMET_ENABLED.key -> "false") { + val expected = sql( + s"SELECT id, $catalog.system.years(y), $catalog.system.months(m), " + + s"$catalog.system.days(d), $catalog.system.hours(h) FROM $table").collect() + checkAnswer( + sql( + "SELECT id, _partition.y_year, _partition.m_month, _partition.d_day, " + + s"_partition.h_hour FROM $table"), + expected) + } + } finally { + sql(s"DROP TABLE IF EXISTS $table") + } + } + } + } + test("non-literal or non-positive parameters fall back to Spark") { withSourceTable { checkSparkAnswerAndFallbackReason( @@ -421,13 +477,50 @@ class CometIcebergSystemFunctionSuite } } - /** Runs `f` with the Iceberg catalog registered and the source parquet table in scope. */ - private def withSourceTable(f: => Unit): Unit = withTempIcebergDir { warehouseDir => + /** Runs `f` with the Iceberg catalog registered. */ + private def withIcebergCatalog(f: => Unit): Unit = withTempIcebergDir { warehouseDir => withSQLConf( s"spark.sql.catalog.$catalog" -> "org.apache.iceberg.spark.SparkCatalog", s"spark.sql.catalog.$catalog.type" -> "hadoop", - s"spark.sql.catalog.$catalog.warehouse" -> warehouseDir.getAbsolutePath) { - withParquetTable(sourcePath, source)(f) + s"spark.sql.catalog.$catalog.warehouse" -> warehouseDir.getAbsolutePath)(f) + } + + /** Runs `f` with the Iceberg catalog registered and the source parquet table in scope. */ + private def withSourceTable(f: => Unit): Unit = + withIcebergCatalog(withParquetTable(sourcePath, source)(f)) + + /** + * Timestamps just after pre-1970 unit boundaries, where Iceberg does not floor. Its + * `DateTimeUtil` places a pre-1970 timestamp whose microsecond of second is 999999 by the + * second before it, so right after a boundary it gets the unit before: 1969-01-01 + * 00:00:00.999999 is in year -2, month -13, day 1968-12-31, and hour -8761, where a floor gives + * -1, -12, 1969-01-01, and -8760. + */ + private val preEpochTimestamps = Seq( + "1969-01-01 00:00:00.999999", // a year, month, day, and hour boundary + "1969-12-01 00:00:00.999999", // a month, day, and hour boundary + "1969-12-31 00:00:00.999999", // a day and hour boundary + "1969-12-31 23:00:00.999999", // an hour boundary + "1969-12-31 22:30:00", + "1968-12-31 12:00:00", + // After the epoch, where Iceberg floors. + "1970-01-01 01:00:00.999999") + + /** + * Runs `f` with a parquet table `pre_epoch (id, ts, ntz)` holding `preEpochTimestamps` as both + * timestamp types. The session timezone is UTC, so that the `TIMESTAMP` values sit on the same + * boundaries as the `TIMESTAMP_NTZ` ones. + */ + private def withPreEpochTable(f: => Unit): Unit = withSQLConf( + SQLConf.SESSION_LOCAL_TIMEZONE.key -> "UTC", + SQLConf.PARQUET_OUTPUT_TIMESTAMP_TYPE.key -> "TIMESTAMP_MICROS") { + withTable("pre_epoch") { + sql("CREATE TABLE pre_epoch (id INT, ts TIMESTAMP, ntz TIMESTAMP_NTZ) USING parquet") + val rows = preEpochTimestamps.zipWithIndex.map { case (timestamp, i) => + s"(${i + 1}, TIMESTAMP '$timestamp', TIMESTAMP_NTZ '$timestamp')" + } + sql(s"INSERT INTO pre_epoch VALUES ${rows.mkString(", ")}") + f } } From 60037f57131da20136a8c5f9b1d44366a0f53119 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 14:23:01 -0600 Subject: [PATCH 2/6] test: make the low-end Iceberg rounding row reach the pre-epoch adjustment i64::MIN + 775_807 ends in 999999 but lies 71_945_999_999 micros into its day and 3_545_999_999 into its hour, so a plain floor gave the same four values as the i64::MIN row. Use -290307-01-01T00:00:00.999999 instead, 999999 micros into the lowest year boundary in range, where all four transforms take the unit before. The expected values come from DateTimeUtil and Transforms in Iceberg 1.5.2, 1.8.1, 1.10.0 and 1.11.0. --- native/spark-expr/src/iceberg_funcs/temporal.rs | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/native/spark-expr/src/iceberg_funcs/temporal.rs b/native/spark-expr/src/iceberg_funcs/temporal.rs index 717e72f48c6..ebcd6b26b55 100644 --- a/native/spark-expr/src/iceberg_funcs/temporal.rs +++ b/native/spark-expr/src/iceberg_funcs/temporal.rs @@ -359,14 +359,15 @@ mod tests { // `(int)` narrowing wraps it exactly as `as i32` does. (i64::MAX, 292_277, 3_507_324, 106_751_991, -1_732_919_508), (i64::MIN, -292_278, -3_507_325, -106_751_992, 1_732_919_507), - // The lowest `i64` that ends in 999999, where moving the value a second earlier would - // overflow. + // -290307-01-01T00:00:00.999999, 999999 micros into the lowest year boundary in + // range, so all four take the unit before it. A floor gives -292_277, -3_507_324, + // -106_751_981, and 1_732_919_752. ( - i64::MIN + 775_807, + -106_751_981 * MICROS_PER_DAY + 999_999, -292_278, -3_507_325, - -106_751_992, - 1_732_919_507, + -106_751_982, + 1_732_919_751, ), ( 8_000_000_000_000_000_000, From d711633555192f1c3a3e004d4f486d0bf9813ddf Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 14:23:08 -0600 Subject: [PATCH 3/6] docs: update comments that still say the Iceberg writer computes only year and month natively The native writer now computes year, month, day and hour of timestamps with Comet's kernels, but a few comments still described the old split: the doc on SparkIcebergTemporalTransform::transform, the note on iceberg_rust_years_follow_the_timezone_tag that days and hours could be delegated to iceberg-rust, the reasons PartitionSplitter replaces iceberg-rust's splitter, and the test comment that delegating year and month would become an option again once iceberg-rust covers the whole calendar. --- .../operators/iceberg_partition_value.rs | 7 +++--- .../src/execution/operators/iceberg_write.rs | 24 +++++++++++-------- .../spark-expr/src/iceberg_funcs/temporal.rs | 6 +++-- 3 files changed, 22 insertions(+), 15 deletions(-) diff --git a/native/core/src/execution/operators/iceberg_partition_value.rs b/native/core/src/execution/operators/iceberg_partition_value.rs index 31b724a3e89..d47ae30fb36 100644 --- a/native/core/src/execution/operators/iceberg_partition_value.rs +++ b/native/core/src/execution/operators/iceberg_partition_value.rs @@ -331,9 +331,10 @@ mod tests { source.data_type() ); } - // The reason Comet computes these two: iceberg-rust's transforms turn every value above - // into a NULL partition value. If this starts failing, iceberg-rust has learned the whole - // domain and delegating becomes an option again. + // The reason Comet computes these two for dates, and one reason for timestamps (see + // `pre_epoch_timestamps_partition_like_iceberg_java` for the other): iceberg-rust's + // transforms turn every value above into a NULL partition value. If this starts failing, + // iceberg-rust has learned the whole domain and delegating dates becomes an option again. for (column, (transform, source, _)) in columns(iceberg_rust.calculate(&batch).unwrap()) .iter() .zip(&cases) diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index baef664096f..52138702d6f 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -855,8 +855,9 @@ fn file_name_prefix(partition_id: i32, task_attempt_id: i64, operation_id: &str) /// /// This replaces iceberg-rust's `RecordBatchPartitionSplitter`, which computes the values with /// iceberg-rust's own transforms -- they turn a `year` or `month` past `chrono`'s calendar into a -/// NULL (apache/datafusion-comet#6145) -- and groups rows through a HashMap, which emits parts in -/// unspecified order. +/// NULL (apache/datafusion-comet#6145) and put some pre-epoch timestamps in a different time +/// partition from iceberg-java (apache/datafusion-comet#6426) -- and groups rows through a +/// HashMap, which emits parts in unspecified order. struct PartitionSplitter { calculator: PartitionValueCalculator, partition_spec: PartitionSpecRef, @@ -3651,14 +3652,17 @@ mod iceberg_rust_transform_parity { } } - /// Why `years` and `months` are not delegated to iceberg-rust even though `bucket`, `days`, - /// and `hours` could be: its kernels go through Arrow's `date_part`, which honours the - /// array's timezone tag, while Iceberg's Java `DateTimeUtil` is always UTC. Comet only ever - /// produces `UTC` and untagged timestamps today, so the parity above holds; this pins the - /// reason the local kernel exists. Reported as apache/iceberg-rust#3142; if this ever fails, - /// iceberg-rust dropped the tag dependency. Delegating is still unsafe until it also covers - /// dates past `chrono`'s calendar, which the writer's own partition values depend on too: - /// see `years_and_months_past_chronos_calendar_match_iceberg_java` (`iceberg_partition_value`). + /// Why `years` and `months` are not delegated to iceberg-rust even though `bucket` could be: + /// its kernels go through Arrow's `date_part`, which honours the array's timezone tag, while + /// Iceberg's Java `DateTimeUtil` is always UTC. Comet only ever produces `UTC` and untagged + /// timestamps today, so the parity above holds; this pins one reason the local kernel exists. + /// Reported as apache/iceberg-rust#3142; if this ever fails, iceberg-rust dropped the tag + /// dependency. Delegating is still unsafe until it also covers dates past `chrono`'s calendar, + /// which the writer's own partition values depend on too: see + /// `years_and_months_past_chronos_calendar_match_iceberg_java` (`iceberg_partition_value`). + /// Nor can `days` and `hours` be delegated: all four time transforms place some pre-epoch + /// timestamps differently from iceberg-java, see + /// `pre_epoch_timestamps_partition_like_iceberg_java` there. #[test] fn iceberg_rust_years_follow_the_timezone_tag() { // 1969-12-31T23:59:59.999999Z, which is 1970-01-01T05:44:59.999999 in Kathmandu. diff --git a/native/spark-expr/src/iceberg_funcs/temporal.rs b/native/spark-expr/src/iceberg_funcs/temporal.rs index ebcd6b26b55..962900b1b9b 100644 --- a/native/spark-expr/src/iceberg_funcs/temporal.rs +++ b/native/spark-expr/src/iceberg_funcs/temporal.rs @@ -217,8 +217,10 @@ impl SparkIcebergTemporalTransform { } /// Applies the transform to a whole column, for callers outside DataFusion's function - /// machinery. The native Iceberg writer computes its `year` and `month` partition values with - /// it, because iceberg-rust's own transforms return NULL past `chrono`'s range. + /// machinery. The native Iceberg writer computes its time partition values of dates and + /// timestamps with it (all but `day` of a date), because iceberg-rust's own transforms return + /// NULL past `chrono`'s range and place some pre-epoch timestamps differently from + /// iceberg-java. pub fn transform(&self, array: &ArrayRef) -> Result { apply_to_array(array, |array| transform_array(self.unit, array)) } From ccbe740e2c0002881d3717d9c9805490122976ae Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 14:38:54 -0600 Subject: [PATCH 4/6] test: move the pre-1970 Iceberg temporal projection and filter checks to a SQL file They read a plain parquet table, so sql-tests/iceberg/temporal_functions_pre_epoch.sql can register the Iceberg catalog with Config lines and run them through CometSqlFileTestSuite. That suite runs on the same CI profiles as CometIcebergSystemFunctionSuite and skips Spark 4.2 the same way. The file adds a NULL row and a filter for each of the four functions. It also uses expect_native(staticinvoke) to pin that Comet evaluates the calls with its own kernels. A call sent through the codegen dispatcher would run Iceberg's Java code and match Spark whatever the kernels return. --- .../iceberg/temporal_functions_pre_epoch.sql | 80 +++++++++++++++++++ .../CometIcebergSystemFunctionSuite.scala | 18 ----- 2 files changed, 80 insertions(+), 18 deletions(-) create mode 100644 spark/src/test/resources/sql-tests/iceberg/temporal_functions_pre_epoch.sql diff --git a/spark/src/test/resources/sql-tests/iceberg/temporal_functions_pre_epoch.sql b/spark/src/test/resources/sql-tests/iceberg/temporal_functions_pre_epoch.sql new file mode 100644 index 00000000000..9dbca435e2c --- /dev/null +++ b/spark/src/test/resources/sql-tests/iceberg/temporal_functions_pre_epoch.sql @@ -0,0 +1,80 @@ +-- 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. + +-- Iceberg's years, months, days, and hours system functions on timestamps just after a pre-1970 +-- unit boundary (https://github.com/apache/datafusion-comet/issues/6426). Iceberg's DateTimeUtil +-- places a pre-1970 timestamp whose microsecond of second is 999999 by the second before it, so +-- right after a boundary it gets the unit before: 1969-01-01 00:00:00.999999 is in year -2, +-- month -13, day 1968-12-31, and hour -8761, where a floor gives -1, -12, 1969-01-01, and -8760. +-- Each query also runs with Comet off, so the expected values come from Iceberg's own functions. +-- CometIcebergSystemFunctionSuite writes the same timestamps with the native Iceberg writer. + +-- No Iceberg spark-runtime is published for Spark 4.2 yet, and the 4.0 runtime the build reuses +-- is binary-incompatible with it, which is also why CometIcebergTestBase reports Iceberg as +-- unavailable there. See https://github.com/apache/datafusion-comet/issues/4969. +-- MaxSparkVersion: 4.1 + +-- Config: spark.sql.catalog.test_cat=org.apache.iceberg.spark.SparkCatalog +-- Config: spark.sql.catalog.test_cat.type=hadoop +-- Config: spark.sql.catalog.test_cat.warehouse=/tmp/comet-iceberg-sql-test +-- UTC, so that the TIMESTAMP values sit on the same boundaries as the TIMESTAMP_NTZ ones. +-- Config: spark.sql.session.timeZone=UTC +-- Config: spark.sql.parquet.outputTimestampType=TIMESTAMP_MICROS + +-- A parquet table rather than an Iceberg one, so that no scan absorbs a filter. Rows 1 to 4 lie +-- 999999 microseconds after a year, a month, a day, and an hour boundary (a year boundary is also +-- a month, day, and hour boundary, and so on). Row 5 sits inside the hour Iceberg gives row 4, +-- and row 6 inside the day it gives row 1, both away from any boundary. Row 7 is after the epoch, +-- where Iceberg floors. +statement +CREATE TABLE iceberg_pre_epoch (id INT, ts TIMESTAMP, ntz TIMESTAMP_NTZ) USING parquet + +statement +INSERT INTO iceberg_pre_epoch VALUES + (1, TIMESTAMP '1969-01-01 00:00:00.999999', TIMESTAMP_NTZ '1969-01-01 00:00:00.999999'), + (2, TIMESTAMP '1969-12-01 00:00:00.999999', TIMESTAMP_NTZ '1969-12-01 00:00:00.999999'), + (3, TIMESTAMP '1969-12-31 00:00:00.999999', TIMESTAMP_NTZ '1969-12-31 00:00:00.999999'), + (4, TIMESTAMP '1969-12-31 23:00:00.999999', TIMESTAMP_NTZ '1969-12-31 23:00:00.999999'), + (5, TIMESTAMP '1969-12-31 22:30:00', TIMESTAMP_NTZ '1969-12-31 22:30:00'), + (6, TIMESTAMP '1968-12-31 12:00:00', TIMESTAMP_NTZ '1968-12-31 12:00:00'), + (7, TIMESTAMP '1970-01-01 01:00:00.999999', TIMESTAMP_NTZ '1970-01-01 01:00:00.999999'), + (8, NULL, NULL) + +-- Every query names staticinvoke, the node Spark plans an Iceberg system function call as, in +-- expect_native. A call Comet did not lower to its own kernel would run Iceberg's Java code through +-- the codegen dispatcher and match Spark whatever the kernel returns. +query expect_native(staticinvoke) +SELECT id, + test_cat.system.years(ts), test_cat.system.months(ts), + test_cat.system.days(ts), test_cat.system.hours(ts), + test_cat.system.years(ntz), test_cat.system.months(ntz), + test_cat.system.days(ntz), test_cat.system.hours(ntz) +FROM iceberg_pre_epoch + +-- Rows 1 and 6. A floor would leave out row 1 in these three. +query expect_native(staticinvoke) +SELECT id FROM iceberg_pre_epoch WHERE test_cat.system.years(ts) = -2 + +query expect_native(staticinvoke) +SELECT id FROM iceberg_pre_epoch WHERE test_cat.system.months(ts) = -13 + +query expect_native(staticinvoke) +SELECT id FROM iceberg_pre_epoch WHERE test_cat.system.days(ts) = DATE '1968-12-31' + +-- Rows 4 and 5. A floor would leave out row 4. +query expect_native(staticinvoke) +SELECT id FROM iceberg_pre_epoch WHERE test_cat.system.hours(ts) = -2 diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala index bf25e23e188..03256b5598a 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala @@ -174,24 +174,6 @@ class CometIcebergSystemFunctionSuite } } - test("years, months, days, and hours match Iceberg on pre-1970 timestamps ending in .999999") { - // The random corpus almost never lands on such a value, and its boundary rows sit a - // microsecond before a boundary rather than after it. - withIcebergCatalog { - withPreEpochTable { - val transformed = for { - column <- Seq("ts", "ntz") - function <- Seq("years", "months", "days", "hours") - } yield s"$catalog.system.$function($column)" - checkSparkAnswerAndOperator(s"SELECT id, ${transformed.mkString(", ")} FROM pre_epoch") - checkSparkAnswerAndOperator( - s"SELECT id FROM pre_epoch WHERE $catalog.system.days(ts) = DATE '1968-12-31'") - checkSparkAnswerAndOperator( - s"SELECT id FROM pre_epoch WHERE $catalog.system.hours(ts) = -2") - } - } - } - test("system functions in filters stay native") { withSourceTable { checkSparkAnswerAndOperator( From c5d00f76c83acb82ba67cb35edf7ce2a44c3d2f4 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 14:39:02 -0600 Subject: [PATCH 5/6] test: share the pre-epoch Iceberg fixture through CometIcebergTestBase CometIcebergSystemFunctionExtensionsSuite repeated the catalog setup and its own copy of the pre_epoch table. The Hadoop catalog, the pre-1970 timestamps and the parquet table now live in CometIcebergTestBase, which both suites use. The trait takes a CometTestBase self-type for withSQLConf, withTable and sql, and every suite that mixes it in already extends CometTestBase. The catalog helper is called withHadoopCatalog so that it does not overload CometIcebergWriteActionSuite's private withIcebergCatalog. The extensions suite's four predicates now run on the shared rows. The table drops the TIMESTAMP_NTZ column, which only the projection test that moved to the SQL file read. --- ...IcebergSystemFunctionExtensionsSuite.scala | 50 +++++++---------- .../CometIcebergSystemFunctionSuite.scala | 50 ++--------------- .../apache/comet/CometIcebergTestBase.scala | 53 +++++++++++++++++-- 3 files changed, 72 insertions(+), 81 deletions(-) diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala index 09b563e929b..e71acf4e5c9 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala @@ -23,7 +23,6 @@ import org.apache.spark.SparkConf import org.apache.spark.sql.CometTestBase import org.apache.spark.sql.catalyst.expressions.ApplyFunctionExpression import org.apache.spark.sql.catalyst.plans.logical.Filter -import org.apache.spark.sql.internal.SQLConf /** * Iceberg's system functions once Iceberg's SQL extensions have rewritten them. @@ -43,37 +42,24 @@ class CometIcebergSystemFunctionExtensionsSuite extends CometTestBase with Comet test("rewritten temporal filters match Iceberg on pre-1970 timestamps ending in .999999") { assume(icebergAvailable, "Iceberg not available in classpath") - withTempIcebergDir { warehouseDir => - withSQLConf( - "spark.sql.catalog.ice" -> "org.apache.iceberg.spark.SparkCatalog", - "spark.sql.catalog.ice.type" -> "hadoop", - "spark.sql.catalog.ice.warehouse" -> warehouseDir.getAbsolutePath, - SQLConf.SESSION_LOCAL_TIMEZONE.key -> "UTC", - SQLConf.PARQUET_OUTPUT_TIMESTAMP_TYPE.key -> "TIMESTAMP_MICROS") { - withTable("pre_epoch") { - // A parquet table, so that no scan absorbs the filter and Comet evaluates it. Iceberg - // places a pre-1970 timestamp ending in .999999 by the second before it: the first row - // is in hour -8761, day 1968-12-31, month -13, and year -2, and the second in hour -2. - sql("CREATE TABLE pre_epoch (id INT, ts TIMESTAMP) USING parquet") - sql("""INSERT INTO pre_epoch VALUES - (1, TIMESTAMP '1969-01-01 00:00:00.999999'), - (2, TIMESTAMP '1969-12-31 23:00:00.999999'), - (3, TIMESTAMP '1969-12-31 22:30:00'), - (4, TIMESTAMP '1968-12-31 12:00:00')""") - Seq( - "ice.system.hours(ts) = -2", - "ice.system.days(ts) = DATE '1968-12-31'", - "ice.system.months(ts) = -13", - "ice.system.years(ts) = -2").foreach { predicate => - val df = sql(s"SELECT id FROM pre_epoch WHERE $predicate") - val plan = df.queryExecution.optimizedPlan - assert( - plan - .collect { case filter: Filter => filter.condition } - .exists(_.find(_.isInstanceOf[ApplyFunctionExpression]).isDefined), - s"expected Iceberg's extensions to rewrite $predicate:\n$plan") - checkSparkAnswerAndOperator(df) - } + withHadoopCatalog("ice") { + withPreEpochTable { + // Iceberg places a pre-1970 timestamp ending in .999999 by the second before it: row 1 is + // in hour -8761, day 1968-12-31, month -13, and year -2, and row 4 in hour -2. So each + // predicate matches one of them besides a row away from any boundary (6, or 5 for hours). + Seq( + "ice.system.hours(ts) = -2", + "ice.system.days(ts) = DATE '1968-12-31'", + "ice.system.months(ts) = -13", + "ice.system.years(ts) = -2").foreach { predicate => + val df = sql(s"SELECT id FROM pre_epoch WHERE $predicate") + val plan = df.queryExecution.optimizedPlan + assert( + plan + .collect { case filter: Filter => filter.condition } + .exists(_.find(_.isInstanceOf[ApplyFunctionExpression]).isDefined), + s"expected Iceberg's extensions to rewrite $predicate:\n$plan") + checkSparkAnswerAndOperator(df) } } } diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala index 03256b5598a..5104bfd8db0 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala @@ -267,8 +267,9 @@ class CometIcebergSystemFunctionSuite // The sort in front of the write runs the native kernels and the native writer computes each // row's partition value, so both have to follow Iceberg for the table to hold the partitions // iceberg-java would have written. A spec takes one time transform per source column, hence a - // column per transform. - withIcebergCatalog { + // column per transform. The kernels' results in projections and filters are checked by + // sql-tests/iceberg/temporal_functions_pre_epoch.sql. + withHadoopCatalog(catalog) { withPreEpochTable { val table = s"$catalog.db.pre_epoch_partitions" sql(s""" @@ -459,52 +460,9 @@ class CometIcebergSystemFunctionSuite } } - /** Runs `f` with the Iceberg catalog registered. */ - private def withIcebergCatalog(f: => Unit): Unit = withTempIcebergDir { warehouseDir => - withSQLConf( - s"spark.sql.catalog.$catalog" -> "org.apache.iceberg.spark.SparkCatalog", - s"spark.sql.catalog.$catalog.type" -> "hadoop", - s"spark.sql.catalog.$catalog.warehouse" -> warehouseDir.getAbsolutePath)(f) - } - /** Runs `f` with the Iceberg catalog registered and the source parquet table in scope. */ private def withSourceTable(f: => Unit): Unit = - withIcebergCatalog(withParquetTable(sourcePath, source)(f)) - - /** - * Timestamps just after pre-1970 unit boundaries, where Iceberg does not floor. Its - * `DateTimeUtil` places a pre-1970 timestamp whose microsecond of second is 999999 by the - * second before it, so right after a boundary it gets the unit before: 1969-01-01 - * 00:00:00.999999 is in year -2, month -13, day 1968-12-31, and hour -8761, where a floor gives - * -1, -12, 1969-01-01, and -8760. - */ - private val preEpochTimestamps = Seq( - "1969-01-01 00:00:00.999999", // a year, month, day, and hour boundary - "1969-12-01 00:00:00.999999", // a month, day, and hour boundary - "1969-12-31 00:00:00.999999", // a day and hour boundary - "1969-12-31 23:00:00.999999", // an hour boundary - "1969-12-31 22:30:00", - "1968-12-31 12:00:00", - // After the epoch, where Iceberg floors. - "1970-01-01 01:00:00.999999") - - /** - * Runs `f` with a parquet table `pre_epoch (id, ts, ntz)` holding `preEpochTimestamps` as both - * timestamp types. The session timezone is UTC, so that the `TIMESTAMP` values sit on the same - * boundaries as the `TIMESTAMP_NTZ` ones. - */ - private def withPreEpochTable(f: => Unit): Unit = withSQLConf( - SQLConf.SESSION_LOCAL_TIMEZONE.key -> "UTC", - SQLConf.PARQUET_OUTPUT_TIMESTAMP_TYPE.key -> "TIMESTAMP_MICROS") { - withTable("pre_epoch") { - sql("CREATE TABLE pre_epoch (id INT, ts TIMESTAMP, ntz TIMESTAMP_NTZ) USING parquet") - val rows = preEpochTimestamps.zipWithIndex.map { case (timestamp, i) => - s"(${i + 1}, TIMESTAMP '$timestamp', TIMESTAMP_NTZ '$timestamp')" - } - sql(s"INSERT INTO pre_epoch VALUES ${rows.mkString(", ")}") - f - } - } + withHadoopCatalog(catalog)(withParquetTable(sourcePath, source)(f)) private val sourceSchema = StructType( Seq( diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergTestBase.scala b/spark/src/test/scala/org/apache/comet/CometIcebergTestBase.scala index 8b2c55505c3..317e2370ac8 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergTestBase.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergTestBase.scala @@ -25,18 +25,20 @@ import java.nio.file.Files import scala.collection.mutable import org.apache.spark.CometListenerBusUtils -import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.{CometTestBase, SparkSession} import org.apache.spark.sql.connector.catalog.{Identifier, TableCatalog} import org.apache.spark.sql.execution.{QueryExecution, SparkPlan} +import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.util.QueryExecutionListener import org.apache.comet.CometSparkSessionExtensions.isSpark42Plus import org.apache.comet.iceberg.IcebergReflection /** - * Shared fixtures for Iceberg-backed test suites: classpath probe and per-test temp directory. + * Shared fixtures for Iceberg-backed test suites: classpath probe, per-test temp directory, a + * Hadoop catalog, and a table of pre-1970 timestamps. Mix in alongside `CometTestBase`. */ -trait CometIcebergTestBase { +trait CometIcebergTestBase { this: CometTestBase => // No Iceberg spark-runtime is published for Spark 4.2 yet, so the build reuses the 4.0 runtime. // That jar is binary-incompatible with Spark 4.2, whose `connector.catalog.View` is a class @@ -137,6 +139,51 @@ trait CometIcebergTestBase { file.delete() } + /** Runs `f` with an Iceberg `hadoop` catalog registered as `catalog`, in a temp warehouse. */ + protected def withHadoopCatalog(catalog: String)(f: => Unit): Unit = + withTempIcebergDir { warehouseDir => + withSQLConf( + s"spark.sql.catalog.$catalog" -> "org.apache.iceberg.spark.SparkCatalog", + s"spark.sql.catalog.$catalog.type" -> "hadoop", + s"spark.sql.catalog.$catalog.warehouse" -> warehouseDir.getAbsolutePath)(f) + } + + /** + * Timestamps just after pre-1970 unit boundaries, where Iceberg does not floor. Its + * `DateTimeUtil` places a pre-1970 timestamp whose microsecond of second is 999999 by the + * second before it, so right after a boundary it gets the unit before: 1969-01-01 + * 00:00:00.999999 is in year -2, month -13, day 1968-12-31, and hour -8761, where a floor gives + * -1, -12, 1969-01-01, and -8760. `sql-tests/iceberg/temporal_functions_pre_epoch.sql` runs the + * system functions over the same timestamps in projections and filters. + */ + protected val preEpochTimestamps: Seq[String] = Seq( + "1969-01-01 00:00:00.999999", // a year, month, day, and hour boundary + "1969-12-01 00:00:00.999999", // a month, day, and hour boundary + "1969-12-31 00:00:00.999999", // a day and hour boundary + "1969-12-31 23:00:00.999999", // an hour boundary + "1969-12-31 22:30:00", + "1968-12-31 12:00:00", + // After the epoch, where Iceberg floors. + "1970-01-01 01:00:00.999999") + + /** + * Runs `f` with a parquet table `pre_epoch (id, ts)` holding `preEpochTimestamps`, numbered + * from 1. A parquet table, so that no scan absorbs a filter on `ts` and Comet evaluates it. The + * session timezone is UTC, so that the values sit on the unit boundaries. + */ + protected def withPreEpochTable(f: => Unit): Unit = withSQLConf( + SQLConf.SESSION_LOCAL_TIMEZONE.key -> "UTC", + SQLConf.PARQUET_OUTPUT_TIMESTAMP_TYPE.key -> "TIMESTAMP_MICROS") { + withTable("pre_epoch") { + sql("CREATE TABLE pre_epoch (id INT, ts TIMESTAMP) USING parquet") + val rows = preEpochTimestamps.zipWithIndex.map { case (timestamp, i) => + s"(${i + 1}, TIMESTAMP '$timestamp')" + } + sql(s"INSERT INTO pre_epoch VALUES ${rows.mkString(", ")}") + f + } + } + /** * The executed plan of every query that ran while `action` ran. Queries that failed are * included only when `includeFailures` is set, which is what an action expected to abort needs. From 8e8513aa3061561ed16809830b3b7391404df5d5 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 06:36:24 -0600 Subject: [PATCH 6/6] docs: link the iceberg-rust issue for its truncating `day` apache/iceberg-rust#3315 reports the bug that moves some pre-epoch timestamps into the next day, so the next reader can tell when the special case can go. --- .../src/execution/operators/iceberg_partition_value.rs | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/native/core/src/execution/operators/iceberg_partition_value.rs b/native/core/src/execution/operators/iceberg_partition_value.rs index d47ae30fb36..970297283a7 100644 --- a/native/core/src/execution/operators/iceberg_partition_value.rs +++ b/native/core/src/execution/operators/iceberg_partition_value.rs @@ -51,7 +51,7 @@ use iceberg::{Error, ErrorKind, Result}; /// partitions (apache/datafusion-comet#6426). /// - `day` moves a timestamp from the last second of a day before 1969-12-31 into the next day, /// unless its microsecond of second is 0 or 999999: it takes the whole seconds with a truncating -/// division and the microseconds with a flooring one. +/// division and the microseconds with a flooring one (apache/iceberg-rust#3315). /// /// A partition value that differs from the sort key can also fail a clustered write, which rejects /// a row whose partition it has already closed. Everywhere else the two implementations agree, so @@ -439,8 +439,9 @@ mod tests { let label = format!("{transform} of {}", source.data_type()); assert_eq!(&comet[i], java, "{label}"); // The reason Comet computes these: iceberg-rust floors the first two rows, and its - // `day` moves the last two into 1969-12-31. If this starts failing, iceberg-rust's - // transforms have changed and delegating needs another look. + // `day` moves the last two into 1969-12-31 (apache/iceberg-rust#3315). If this + // starts failing, iceberg-rust's transforms have changed and delegating needs + // another look. assert_eq!(&iceberg_rust[i], rust, "iceberg-rust's {label}"); } }