diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 2aa61ff811b..8befa94b7ca 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -471,6 +471,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 3660aa88842..358632ac991 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..970297283a7 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 (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 +/// 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)?), }) } @@ -312,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) @@ -328,10 +348,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 +383,74 @@ 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 (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}"); + } + } + } + + /// 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 +507,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..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, @@ -3314,11 +3315,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 +3570,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 +3607,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(); @@ -3648,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 96a2ee82d0b..962900b1b9b 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. @@ -191,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)) } @@ -231,6 +259,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 +361,16 @@ 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), + // -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. + ( + -106_751_981 * MICROS_PER_DAY + 999_999, + -292_278, + -3_507_325, + -106_751_982, + 1_732_919_751, + ), ( 8_000_000_000_000_000_000, 253_509, @@ -373,6 +437,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/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/CometIcebergSystemFunctionExtensionsSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala new file mode 100644 index 00000000000..e71acf4e5c9 --- /dev/null +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionExtensionsSuite.scala @@ -0,0 +1,67 @@ +/* + * 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 + +/** + * 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") + 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 118b2371d53..5104bfd8db0 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala @@ -263,6 +263,45 @@ 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. 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""" + 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( @@ -422,14 +461,8 @@ class CometIcebergSystemFunctionSuite } /** Runs `f` with the Iceberg catalog registered and the source parquet table in scope. */ - private def withSourceTable(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) - } - } + private def withSourceTable(f: => Unit): Unit = + 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.