Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
15 changes: 10 additions & 5 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
146 changes: 118 additions & 28 deletions native/core/src/execution/operators/iceberg_partition_value.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question: is there an iceberg-rust issue for the truncating day? I couldn't find one. If not, would it be worth filing one and linking it here, as is done for apache/iceberg-rust#3142 on the year and month tag? That would tell the next reader when this special case can go.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No, there isn't one. apache/iceberg-rust#3022 even says Day already aligns with Java, but day_timestamp_micro still truncates the seconds toward zero. I've written one up with a reproducer and a one-line fix, and I'll link it here once it's filed.

@andygrove andygrove Oct 1, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Filed as apache/iceberg-rust#3315, and linked in 8e8513a, both in the doc comment here and in the test that pins iceberg-rust's values. That test will start failing when Comet picks up a fix, which is the cue to look at delegating day again. year, month and hour would still need the kernels for the 999999 rows.

/// 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<FieldTransform>,
Expand Down Expand Up @@ -125,19 +136,27 @@ enum FieldTransform {

impl FieldTransform {
fn try_new(transform: Transform, source_type: Option<&Type>) -> Result<Self> {
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())
}
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)?),
})
}
Expand Down Expand Up @@ -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)
Expand All @@ -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)),
Expand Down Expand Up @@ -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(&[
Expand Down Expand Up @@ -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::<Vec<_>>();
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())
Expand Down
51 changes: 29 additions & 22 deletions native/core/src/execution/operators/iceberg_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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::{
Expand Down Expand Up @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: a couple of doc comments the PR doesn't touch still describe the old split. The one on iceberg_rust_years_follow_the_timezone_tag below says days and hours could be delegated to iceberg-rust. SparkIcebergTemporalTransform::transform in temporal.rs says the writer uses it for year and month only.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed both in d711633, along with two more that a grep turned up: the PartitionSplitter doc, and a test comment in iceberg_partition_value.rs that still said delegating year and month would become possible again.

/// 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![
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading