Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
41 commits
Select commit Hold shift + click to select a range
72997e9
fix: preserve CalendarInterval microseconds
peterxcli Aug 7, 2026
a0cb4dd
scalafix
peterxcli Aug 7, 2026
c7f906c
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 7, 2026
9b764d8
fix: use Arrow null check for struct vectors
peterxcli Aug 7, 2026
ee5ea96
andy's review
peterxcli Aug 13, 2026
c0f2f9b
Update make_interval_ansi.sql
peterxcli Aug 14, 2026
7418d88
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 14, 2026
76f39b5
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 18, 2026
bf29348
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 19, 2026
ae39d41
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 22, 2026
112f5c5
Merge branch 'main' into feat/full-native-make-interval
peterxcli Aug 24, 2026
eb4e1ce
fix scalalint
peterxcli Aug 24, 2026
8b61b8a
docs: explain metadata-marker requirement in isCalendarIntervalStruct…
peterxcli Aug 27, 2026
88ab0b0
Merge remote-tracking branch 'upstream/main' into feat/full-native-ma…
peterxcli Aug 27, 2026
a2f8885
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 4, 2026
91f65ae
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 4, 2026
fef4551
fix: keep datafusion-spark for benchmarks
peterxcli Sep 4, 2026
6612649
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 4, 2026
4a229be
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 5, 2026
2857512
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 9, 2026
70874ce
Merge remote-tracking branch 'upstream/main' into feat/full-native-ma…
peterxcli Sep 13, 2026
b55bde0
fix: restore datafusion-spark as a spark-expr library dependency
peterxcli Sep 13, 2026
48b15ac
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 15, 2026
99087b8
fix: update native make_interval routing expectations and lint
peterxcli Sep 15, 2026
fb87be6
Merge remote-tracking branch 'upstream/main' into feat/full-native-ma…
peterxcli Sep 16, 2026
2074ce9
Merge branch 'main' into feat/full-native-make-interval
peterxcli Sep 17, 2026
f84ff98
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Sep 26, 2026
6742caa
fix: hash CalendarInterval values like Spark in the native hash kernels
peterxcli Sep 28, 2026
7cb1f76
fix: keep CalendarInterval layout through collect_list normalization
peterxcli Sep 28, 2026
05fedc4
fix: preserve make_interval NULL short-circuiting for eager arguments
peterxcli Sep 28, 2026
0a90d3d
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Oct 2, 2026
10e26d3
fix: drop the obsolete calendar interval range cap on folded literals
peterxcli Oct 2, 2026
c017d8f
test: cover calendar interval negation, addition and subtraction
peterxcli Oct 2, 2026
5c9b6a5
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Oct 2, 2026
039ebe5
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Oct 4, 2026
603e7a4
fix: pass Decimal(18, 6) seconds in the make_interval bench
peterxcli Oct 4, 2026
6cd802a
fix: dispatch make_interval when an argument has no native path
peterxcli Oct 4, 2026
7c043d1
fix: keep Python operators with calendar intervals on Spark
peterxcli Oct 4, 2026
b0abd19
perf: write calendar interval child vectors directly in the Arrow writer
peterxcli Oct 4, 2026
4f4b8f1
test: cover more consumers of native make_interval output
peterxcli Oct 4, 2026
d0f2923
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Oct 4, 2026
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
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/expressions.md
Original file line number Diff line number Diff line change
Expand Up @@ -282,7 +282,7 @@ The type-name conversion functions (`bigint`, `binary`, `boolean`, `date`, `deci
| `localtimestamp` | ✅ | — | |
| `make_date` | ✅ | Native | |
| `make_dt_interval` | ✅ | Codegen dispatch | |
| `make_interval` | ✅ | Hybrid | Routes through the JVM codegen dispatcher by default; intervals outside Arrow's nanosecond range are tracked by [#5279](https://github.com/apache/datafusion-comet/issues/5279); the native path is opt-in via allowIncompatible ([details](compatibility/expressions/datetime.md)) |
| `make_interval` | ✅ | Hybrid | Runs natively when each argument after the first nullable one is a column, a literal, or a lossless up-cast of one. Other shapes use the JVM codegen dispatcher to keep Spark's NULL short-circuit. An argument without a native path also uses the dispatcher ([details](compatibility/expressions/datetime.md)) |
| `make_time` | ✅ | — | Spark 4.1+; requires `spark.sql.timeType.enabled=true`, which Spark leaves off by default. Runs natively; remaining TIME type work is tracked by [#4288](https://github.com/apache/datafusion-comet/issues/4288) |
| `make_timestamp` | ✅ | Hybrid | |
| `make_timestamp_ltz` | ✅ | — | 2-arg TIME form falls back |
Expand Down
56 changes: 52 additions & 4 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,10 +85,11 @@ use datafusion_comet_operators::{
WindowFnKind,
};
use datafusion_comet_spark_expr::{
create_comet_physical_fun, create_comet_physical_fun_with_eval_mode, BinaryOutputStyle,
BloomFilterAgg, BloomFilterMightContain, CheckedBinaryExpr, CometCollectList, CometCollectSet,
CsvWriteOptions, EvalMode, ListPositionsExpr, SparkArraysZipFunc, SparkBloomFilterVersion,
SparkListAgg, SparkPercentile, Subquery, SumInteger, ToCsv,
calendar_interval_type, create_comet_physical_fun, create_comet_physical_fun_with_eval_mode,
is_calendar_interval_fields, BinaryOutputStyle, BloomFilterAgg, BloomFilterMightContain,
CheckedBinaryExpr, CometCollectList, CometCollectSet, CsvWriteOptions, EvalMode,
ListPositionsExpr, SparkArraysZipFunc, SparkBloomFilterVersion, SparkListAgg, SparkPercentile,
Subquery, SumInteger, ToCsv,
};
use datafusion_datasource::TableSchema;
use iceberg::expr::Bind;
Expand Down Expand Up @@ -170,6 +171,13 @@ struct JoinParameters {

/// Return a copy of `data_type` with every nested field marked nullable. Map key fields are left
/// non-nullable to preserve Arrow's map invariant. Primitive types are returned unchanged.
///
/// A Spark `CalendarIntervalType` value, carried as the tagged struct of
/// [`calendar_interval_type`], is a scalar with a fixed layout rather than a struct, so it is kept
/// in exactly that layout wherever it appears. Widening its children would make a global
/// `collect_list` declare a result type that differs from the `ARRAY<INTERVAL>` state its final
/// stage reads back in the Spark-declared layout, and the final aggregate would then fail
/// validating its own output batch.
fn make_all_fields_nullable(data_type: &DataType) -> DataType {
fn nullable_field(field: &Field, nullable: bool) -> FieldRef {
Arc::new(
Expand All @@ -182,6 +190,7 @@ fn make_all_fields_nullable(data_type: &DataType) -> DataType {
)
}
match data_type {
DataType::Struct(fields) if is_calendar_interval_fields(fields) => calendar_interval_type(),
DataType::Struct(fields) => {
DataType::Struct(fields.iter().map(|f| nullable_field(f, true)).collect())
}
Expand Down Expand Up @@ -6887,6 +6896,45 @@ mod tests {
);
}

/// A `CalendarIntervalType` value keeps its tagged layout through the coercion, bare or
/// nested, while ordinary fields around it are still widened. A global `collect_list` reads
/// its state back in the Spark-declared `ARRAY<INTERVAL>` layout in the final stage, so a
/// widened declaration there would fail the aggregate's output-batch validation.
#[test]
fn test_collect_agg_coercion_keeps_calendar_interval_layout() {
let interval = datafusion_comet_spark_expr::calendar_interval_type();
let raw = Arc::new(Column::new("s", 0)) as Arc<dyn PhysicalExpr>;
let coerced_type = |data_type: DataType| {
let plan_schema = collect_agg_schema(data_type);
PhysicalPlanner::coerce_collect_child_nullability(Arc::clone(&raw), &plan_schema)
.unwrap()
.data_type(plan_schema.as_ref())
.unwrap()
};

assert_eq!(coerced_type(interval.clone()), interval);

let nested = |nullable: bool| {
DataType::Struct(Fields::from(vec![
Field::new("i", interval.clone(), true),
Field::new("n", DataType::Int32, nullable),
]))
};
assert_eq!(coerced_type(nested(false)), nested(true));

// An interval whose children were already widened is normalized back to the layout.
let DataType::Struct(fields) = &interval else {
unreachable!()
};
let widened = DataType::Struct(
fields
.iter()
.map(|f| Arc::new(f.as_ref().clone().with_nullable(true)))
.collect(),
);
assert_eq!(coerced_type(widened), interval);
}

/// Primitive arguments carry no field-level nullability, so no cast is inserted.
#[test]
fn test_collect_agg_no_coercion_for_primitive_child() {
Expand Down
5 changes: 2 additions & 3 deletions native/core/src/execution/serde.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ use datafusion_comet_proto::{
spark_expression::DataType,
spark_operator,
};
use datafusion_comet_spark_expr::calendar_interval_type;
use parquet::{arrow::PARQUET_FIELD_ID_META_KEY, variant::VariantType};
use prost::Message;
use std::{io::Cursor, sync::Arc};
Expand Down Expand Up @@ -103,9 +104,7 @@ pub fn to_arrow_datatype(dt_value: &DataType) -> ArrowDataType {
// Spark's DayTimeIntervalType stores microseconds in an int64, which matches Arrow
// Duration(Microsecond) rather than the lossy Interval(DayTime) {days, millis} layout.
DataTypeId::DayTimeInterval => ArrowDataType::Duration(TimeUnit::Microsecond),
// Spark's CalendarIntervalType stores months, days, and microseconds. Arrow stores the
// same components with nanosecond precision.
DataTypeId::CalendarInterval => ArrowDataType::Interval(IntervalUnit::MonthDayNano),
DataTypeId::CalendarInterval => calendar_interval_type(),
DataTypeId::Variant => ArrowDataType::Struct(Fields::from(vec![
Field::new("value", ArrowDataType::Binary, false),
Field::new("metadata", ArrowDataType::Binary, false),
Expand Down
18 changes: 15 additions & 3 deletions native/spark-expr/benches/make_interval.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
// specific language governing permissions and limitations
// under the License.

use arrow::array::{ArrayRef, Decimal128Array};
use arrow::datatypes::{DataType, Field};
use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion};
use datafusion::common::config::ConfigOptions;
Expand All @@ -25,22 +26,33 @@ use std::sync::Arc;

#[path = "common/mod.rs"]
mod common;
use common::{f64_array, i32_array, NULL_RATIOS, ROW_COUNTS};
use common::{i32_array, is_null, NULL_RATIOS, ROW_COUNTS};

/// The kernel accepts only `Decimal(18, 6)` seconds, the type Spark gives the `secs` argument.
fn seconds_array(rows: usize, null_ratio: f64) -> ArrayRef {
let arr = (0..rows)
.map(|i| (!is_null(i, null_ratio)).then_some((i % 60_000_000) as i128))
.collect::<Decimal128Array>()
.with_precision_and_scale(18, 6)
.unwrap();
Arc::new(arr)
}

fn criterion_benchmark(c: &mut Criterion) {
let udf = SparkMakeInterval::new(false);
let mut group = c.benchmark_group("make_interval");
for rows in ROW_COUNTS {
for (null_ratio, tag) in NULL_RATIOS {
// make_interval(years, months, weeks, days, hours, mins: Int32, secs: Float64)
// make_interval(years, months, weeks, days, hours, mins: Int32,
// secs: Decimal128(18, 6))
let args = vec![
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 10) as i32)),
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 12) as i32)),
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 4) as i32)),
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 28) as i32)),
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 24) as i32)),
ColumnarValue::Array(i32_array(rows, null_ratio, |i| (i % 60) as i32)),
ColumnarValue::Array(f64_array(rows, null_ratio, |i| (i % 60) as f64)),
ColumnarValue::Array(seconds_array(rows, null_ratio)),
];
group.bench_with_input(
BenchmarkId::from_parameter(format!("{rows}/{tag}")),
Expand Down
Loading
Loading