From fcdab25c3a73f35d5924d0742c5c27326ce7163b Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 20 Jul 2026 21:33:37 +0800 Subject: [PATCH 1/4] feat: support interval codegen dispatch --- .../CometBatchKernelCodegenInput.scala | 30 ++++++++++++------- .../CometBatchKernelCodegenOutput.scala | 5 ++-- .../CometSpecializedGettersDispatch.scala | 5 ++-- .../comet/serde/operator/CometSink.scala | 17 ++++++++--- .../udf/codegen/CometScalaUDFCodegen.scala | 2 +- .../sql/comet/CometLocalTableScanExec.scala | 9 +++--- .../shuffle/CometShuffleExchangeExec.scala | 5 ++-- .../expressions/datetime/make_dt_interval.sql | 16 ++++++++++ .../expressions/datetime/make_ym_interval.sql | 16 ++++++++++ 9 files changed, 78 insertions(+), 27 deletions(-) diff --git a/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenInput.scala b/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenInput.scala index f31c712c62..aac9d5b99b 100644 --- a/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenInput.scala +++ b/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenInput.scala @@ -64,7 +64,9 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { classOf[DateDayVector], classOf[TimeNanoVector], classOf[TimeStampMicroVector], - classOf[TimeStampMicroTZVector]) + classOf[TimeStampMicroTZVector], + classOf[IntervalYearVector], + classOf[DurationVector]) private val cometPlainVectorName: String = classOf[CometPlainVector].getName /** Emit kernel typed-vector field declarations for every level of every input column. */ @@ -129,7 +131,8 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { } val intCases = withOrd.collect { case (ArrowColumnSpec(cls, _), ord) - if cls == classOf[IntVector] || cls == classOf[DateDayVector] => + if cls == classOf[IntVector] || cls == classOf[DateDayVector] || + cls == classOf[IntervalYearVector] => s" case $ord: return this.col$ord.getInt(this.rowIdx);" } val longCases = withOrd.collect { @@ -137,7 +140,8 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { if cls == classOf[BigIntVector] || cls == classOf[TimeNanoVector] || cls == classOf[TimeStampMicroVector] || - cls == classOf[TimeStampMicroTZVector] => + cls == classOf[TimeStampMicroTZVector] || + cls == classOf[DurationVector] => s" case $ord: return this.col$ord.getLong(this.rowIdx);" } val floatCases = withOrd.collect { @@ -593,8 +597,9 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { case BooleanType => s"getBoolean($idx)" case ByteType => s"getByte($idx)" case ShortType => s"getShort($idx)" - case IntegerType | DateType => s"getInt($idx)" - case LongType | TimestampType | TimestampNTZType => s"getLong($idx)" + case IntegerType | DateType | _: YearMonthIntervalType => s"getInt($idx)" + case LongType | TimestampType | TimestampNTZType | _: DayTimeIntervalType => + s"getLong($idx)" case dt if isTimeType(dt) => s"getLong($idx)" case FloatType => s"getFloat($idx)" case DoubleType => s"getDouble($idx)" @@ -691,12 +696,12 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { | public short getShort(int i) { | return $childField.getShort(startIndex + i); | }""".stripMargin - case IntegerType | DateType => + case IntegerType | DateType | _: YearMonthIntervalType => s""" @Override | public int getInt(int i) { | return $childField.getInt(startIndex + i); | }""".stripMargin - case LongType | TimestampType | TimestampNTZType => + case LongType | TimestampType | TimestampNTZType | _: DayTimeIntervalType => s""" @Override | public long getLong(int i) { | return $childField.getLong(startIndex + i); @@ -852,9 +857,9 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { s" case $fi: return ${path}_f$fi.getByte(this.rowIdx);" case ShortType => s" case $fi: return ${path}_f$fi.getShort(this.rowIdx);" - case IntegerType | DateType => + case IntegerType | DateType | _: YearMonthIntervalType => s" case $fi: return ${path}_f$fi.getInt(this.rowIdx);" - case LongType | TimestampType | TimestampNTZType => + case LongType | TimestampType | TimestampNTZType | _: DayTimeIntervalType => s" case $fi: return ${path}_f$fi.getLong(this.rowIdx);" case dt if isTimeType(dt) => s" case $fi: return ${path}_f$fi.getLong(this.rowIdx);" @@ -902,13 +907,16 @@ private[codegen] object CometBatchKernelCodegenInput extends CometTypeShim { fieldReadScalar(fi, ShortType, f.nullable) } val intCases = scalarOrd.collect { - case (f, fi) if f.sparkType == IntegerType || f.sparkType == DateType => + case (f, fi) + if f.sparkType == IntegerType || f.sparkType == DateType || + f.sparkType.isInstanceOf[YearMonthIntervalType] => fieldReadScalar(fi, IntegerType, f.nullable) } val longCases = scalarOrd.collect { case (f, fi) if f.sparkType == LongType || f.sparkType == TimestampType || - f.sparkType == TimestampNTZType || isTimeType(f.sparkType) => + f.sparkType == TimestampNTZType || isTimeType(f.sparkType) || + f.sparkType.isInstanceOf[DayTimeIntervalType] => fieldReadScalar(fi, LongType, f.nullable) } val floatCases = diff --git a/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala b/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala index 4160c478d1..00c6280543 100644 --- a/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala +++ b/spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala @@ -401,8 +401,9 @@ private[codegen] object CometBatchKernelCodegenOutput extends CometTypeShim { case BooleanType => s"$target.getBoolean($idx)" case ByteType => s"$target.getByte($idx)" case ShortType => s"$target.getShort($idx)" - case IntegerType | DateType => s"$target.getInt($idx)" - case LongType | TimestampType | TimestampNTZType => s"$target.getLong($idx)" + case IntegerType | DateType | _: YearMonthIntervalType => s"$target.getInt($idx)" + case LongType | TimestampType | TimestampNTZType | _: DayTimeIntervalType => + s"$target.getLong($idx)" case dt if isTimeType(dt) => s"$target.getLong($idx)" case FloatType => s"$target.getFloat($idx)" case DoubleType => s"$target.getDouble($idx)" diff --git a/spark/src/main/scala/org/apache/comet/codegen/CometSpecializedGettersDispatch.scala b/spark/src/main/scala/org/apache/comet/codegen/CometSpecializedGettersDispatch.scala index 2f81c58c06..3eeca4e404 100644 --- a/spark/src/main/scala/org/apache/comet/codegen/CometSpecializedGettersDispatch.scala +++ b/spark/src/main/scala/org/apache/comet/codegen/CometSpecializedGettersDispatch.scala @@ -40,8 +40,9 @@ private[codegen] object CometSpecializedGettersDispatch { case BooleanType => java.lang.Boolean.valueOf(g.getBoolean(ordinal)) case ByteType => java.lang.Byte.valueOf(g.getByte(ordinal)) case ShortType => java.lang.Short.valueOf(g.getShort(ordinal)) - case IntegerType | DateType => java.lang.Integer.valueOf(g.getInt(ordinal)) - case LongType | TimestampType | TimestampNTZType => + case IntegerType | DateType | _: YearMonthIntervalType => + java.lang.Integer.valueOf(g.getInt(ordinal)) + case LongType | TimestampType | TimestampNTZType | _: DayTimeIntervalType => java.lang.Long.valueOf(g.getLong(ordinal)) case FloatType => java.lang.Float.valueOf(g.getFloat(ordinal)) case DoubleType => java.lang.Double.valueOf(g.getDouble(ordinal)) diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometSink.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometSink.scala index b1834c5083..812c6903e4 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometSink.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometSink.scala @@ -26,6 +26,7 @@ import org.apache.spark.sql.comet.execution.shuffle.CometShuffleExchangeExec import org.apache.spark.sql.execution.SparkPlan import org.apache.spark.sql.execution.adaptive.ShuffleQueryStageExec import org.apache.spark.sql.execution.exchange.ReusedExchangeExec +import org.apache.spark.sql.types.{ArrayType, DataType, DayTimeIntervalType, MapType, StructType, YearMonthIntervalType} import org.apache.comet.CometConf import org.apache.comet.CometSparkSessionExtensions.withFallbackReason @@ -42,12 +43,21 @@ abstract class CometSink[T <: SparkPlan] extends CometOperatorSerde[T] { override def enabledConfig: Option[ConfigEntry[Boolean]] = None + protected final def supportedSinkDataType(dt: DataType): Boolean = dt match { + case _: YearMonthIntervalType | _: DayTimeIntervalType => true + case StructType(fields) => + fields.nonEmpty && fields.forall(f => supportedSinkDataType(f.dataType)) + case ArrayType(elementType, _) => supportedSinkDataType(elementType) + case MapType(keyType, valueType, _) => + supportedSinkDataType(keyType) && supportedSinkDataType(valueType) + case _ => supportedDataType(dt) + } + override def convert( op: T, builder: Operator.Builder, childOp: OperatorOuterClass.Operator*): Option[OperatorOuterClass.Operator] = { - val supportedTypes = - op.output.forall(a => supportedDataType(a.dataType, allowComplex = true)) + val supportedTypes = op.output.forall(a => supportedSinkDataType(a.dataType)) if (!supportedTypes) { withFallbackReason(op, "Unsupported data type") @@ -115,8 +125,7 @@ object CometExchangeSink extends CometSink[SparkPlan] { private def convertToShuffleScan( op: SparkPlan, builder: Operator.Builder): Option[OperatorOuterClass.Operator] = { - val supportedTypes = - op.output.forall(a => supportedDataType(a.dataType, allowComplex = true)) + val supportedTypes = op.output.forall(a => supportedSinkDataType(a.dataType)) if (!supportedTypes) { withFallbackReason(op, "Unsupported data type for shuffle direct read") diff --git a/spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala b/spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala index f575dd5b53..142a1027df 100644 --- a/spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala +++ b/spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala @@ -220,7 +220,7 @@ class CometScalaUDFCodegen extends CometUDF with Logging { case _: BitVector | _: TinyIntVector | _: SmallIntVector | _: IntVector | _: BigIntVector | _: Float4Vector | _: Float8Vector | _: DecimalVector | _: VarCharVector | _: VarBinaryVector | _: DateDayVector | _: TimeStampMicroVector | - _: TimeStampMicroTZVector => + _: TimeStampMicroTZVector | _: IntervalYearVector | _: DurationVector => ScalarColumnSpec(v.getClass.asInstanceOf[Class[_ <: ValueVector]], nullable = true) case other => throw new UnsupportedOperationException( diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala index 161a4e0553..b4be9679ec 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala @@ -31,7 +31,7 @@ import org.apache.spark.sql.comet.execution.arrow.{CometArrowStream, CometNative import org.apache.spark.sql.comet.util.Utils import org.apache.spark.sql.execution.{LeafExecNode, LocalTableScanExec} import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics} -import org.apache.spark.sql.types.{DataType, NullType} +import org.apache.spark.sql.types.{DataType, DayTimeIntervalType, NullType, YearMonthIntervalType} import com.google.common.base.Objects @@ -118,14 +118,13 @@ object CometLocalTableScanExec extends CometSink[LocalTableScanExec] with DataTy override def enabledConfig: Option[ConfigEntry[Boolean]] = Some( CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED) - // ArrowWriter (used by RowArrowReader) handles NullType via Utils.toArrowType + NullWriter; - // other types off DataTypeSupport's allow list (TimeType, intervals, ...) have no ArrowWriter - // coverage and must fall back to Spark. + // ArrowWriter (used by RowArrowReader) handles these types even though the shared scan allow + // list does not. Keep this local so interval file scans remain unsupported. override def isTypeSupported( dt: DataType, name: String, fallbackReasons: ListBuffer[String]): Boolean = dt match { - case _: NullType => true + case _: NullType | _: YearMonthIntervalType | _: DayTimeIntervalType => true case _ => super.isTypeSupported(dt, name, fallbackReasons) } diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala index c9fe324bd8..71d1784630 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala @@ -41,7 +41,7 @@ import org.apache.spark.sql.execution.adaptive.ShuffleQueryStageExec import org.apache.spark.sql.execution.exchange.{ENSURE_REQUIREMENTS, ShuffleExchangeExec, ShuffleExchangeLike, ShuffleOrigin} import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics, SQLShuffleReadMetricsReporter, SQLShuffleWriteMetricsReporter} import org.apache.spark.sql.internal.SQLConf -import org.apache.spark.sql.types.{ArrayType, BinaryType, BooleanType, ByteType, DataType, DateType, DecimalType, DoubleType, FloatType, IntegerType, LongType, MapType, NullType, ShortType, StringType, StructField, StructType, TimestampNTZType, TimestampType} +import org.apache.spark.sql.types.{ArrayType, BinaryType, BooleanType, ByteType, DataType, DateType, DayTimeIntervalType, DecimalType, DoubleType, FloatType, IntegerType, LongType, MapType, NullType, ShortType, StringType, StructField, StructType, TimestampNTZType, TimestampType, YearMonthIntervalType} import org.apache.spark.sql.vectorized.ColumnarBatch import org.apache.spark.util.MutablePair import org.apache.spark.util.collection.unsafe.sort.{PrefixComparators, RecordComparator} @@ -412,7 +412,8 @@ object CometShuffleExchangeExec def supportedSerializableDataType(dt: DataType): Boolean = dt match { case _: BooleanType | _: ByteType | _: ShortType | _: IntegerType | _: LongType | _: FloatType | _: DoubleType | _: StringType | _: BinaryType | _: TimestampType | - _: TimestampNTZType | _: DecimalType | _: DateType | _: NullType => + _: TimestampNTZType | _: DecimalType | _: DateType | _: NullType | + _: YearMonthIntervalType | _: DayTimeIntervalType => true case dt if isTimeType(dt) => true diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql index 2c61040de9..38e1757dfc 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql @@ -15,6 +15,9 @@ -- specific language governing permissions and limitations -- under the License. +-- Config: spark.comet.exec.localTableScan.enabled=true +-- Config: spark.comet.exec.shuffle.mode=native + -- Routes make_dt_interval through the codegen dispatcher; produces DayTimeIntervalType. statement @@ -34,6 +37,19 @@ SELECT make_dt_interval(1, 2, 3, 4.5), make_dt_interval(0, 0, 0, 0) query SELECT make_dt_interval(1), make_dt_interval(1, 2), make_dt_interval() +-- nested interval output through LocalTableScan and codegen dispatch +query +SELECT transform(a, x -> x) AS result +FROM VALUES + (array(make_dt_interval(1, 2, 3, 4.5), CAST(NULL AS INTERVAL DAY TO SECOND))) +AS t(a) + +-- interval output through native shuffle +query +SELECT d, h, mi, s, make_dt_interval(d, h, mi, s) AS i +FROM test_mdi +DISTRIBUTE BY d + -- overflow: days * MICROS_PER_DAY exceeds the int64 microsecond range. makeDayTimeInterval throws -- unconditionally (not ANSI-gated); this confirms the dispatched codegen path propagates Spark's -- exception. The pattern is the lowercase word so it matches every version: Spark 4.x raises diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql index a642b4925d..4d557b07d9 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql @@ -15,6 +15,9 @@ -- specific language governing permissions and limitations -- under the License. +-- Config: spark.comet.exec.localTableScan.enabled=true +-- Config: spark.comet.exec.shuffle.mode=native + -- Routes make_ym_interval through the codegen dispatcher; produces YearMonthIntervalType. statement @@ -34,6 +37,19 @@ SELECT make_ym_interval(1, 2), make_ym_interval(0, 0), make_ym_interval(-5, 11) query SELECT make_ym_interval(3), make_ym_interval() +-- nested interval output through LocalTableScan and codegen dispatch +query +SELECT transform(a, x -> x) AS result +FROM VALUES + (array(make_ym_interval(1, 2), CAST(NULL AS INTERVAL YEAR TO MONTH))) +AS t(a) + +-- interval output through native shuffle +query +SELECT y, m, make_ym_interval(y, m) AS i +FROM test_myi +DISTRIBUTE BY y + -- overflow: years * 12 exceeds Int range. makeYearMonthInterval throws unconditionally (not -- ANSI-gated); this confirms the dispatched codegen path propagates Spark's exception. The -- pattern is the lowercase word so it matches every version: Spark 4.x raises From 31dc88c3edcb66ce57aa827a3e2bd5be03b826e4 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Fri, 24 Jul 2026 23:48:41 +0800 Subject: [PATCH 2/4] add struct and map with nested interval type --- .../sql-tests/expressions/datetime/make_dt_interval.sql | 7 +++++-- .../sql-tests/expressions/datetime/make_ym_interval.sql | 7 +++++-- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql index 38e1757dfc..c8fab97a80 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql @@ -44,9 +44,12 @@ FROM VALUES (array(make_dt_interval(1, 2, 3, 4.5), CAST(NULL AS INTERVAL DAY TO SECOND))) AS t(a) --- interval output through native shuffle +-- top-level, struct, and map interval output through native shuffle query -SELECT d, h, mi, s, make_dt_interval(d, h, mi, s) AS i +SELECT d, h, mi, s, + make_dt_interval(d, h, mi, s) AS i, + named_struct('i', make_dt_interval(d, h, mi, s)) AS st, + map('i', make_dt_interval(d, h, mi, s)) AS m FROM test_mdi DISTRIBUTE BY d diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql index 4d557b07d9..b60e4efc52 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql @@ -44,9 +44,12 @@ FROM VALUES (array(make_ym_interval(1, 2), CAST(NULL AS INTERVAL YEAR TO MONTH))) AS t(a) --- interval output through native shuffle +-- top-level, struct, and map interval output through native shuffle query -SELECT y, m, make_ym_interval(y, m) AS i +SELECT y, m, + make_ym_interval(y, m) AS i, + named_struct('i', make_ym_interval(y, m)) AS s, + map('i', make_ym_interval(y, m)) AS m FROM test_myi DISTRIBUTE BY y From 12b835a1b965564ab5bfbd3b778fae2df4ebd158 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Fri, 24 Jul 2026 23:51:13 +0800 Subject: [PATCH 3/4] fix scala style --- .../org/apache/spark/sql/comet/CometLocalTableScanExec.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala index 99be4dec0e..d5bd16bdae 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometLocalTableScanExec.scala @@ -31,7 +31,7 @@ import org.apache.spark.sql.comet.execution.arrow.{CometArrowStream, CometNative import org.apache.spark.sql.comet.util.Utils import org.apache.spark.sql.execution.{LeafExecNode, LocalTableScanExec} import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics} -import org.apache.spark.sql.types.{DataType, NullType, StructType, DayTimeIntervalType, YearMonthIntervalType} +import org.apache.spark.sql.types.{DataType, DayTimeIntervalType, NullType, StructType, YearMonthIntervalType} import com.google.common.base.Objects From 79044828fb4815aecd7dc73261ca6adf652d220d Mon Sep 17 00:00:00 2001 From: peterxcli Date: Sat, 25 Jul 2026 22:55:04 +0800 Subject: [PATCH 4/4] test: cover null intervals in native shuffle --- .../expressions/datetime/make_dt_interval.sql | 11 +++++++++++ .../expressions/datetime/make_ym_interval.sql | 11 +++++++++++ 2 files changed, 22 insertions(+) diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql index c8fab97a80..371f01be43 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_dt_interval.sql @@ -53,6 +53,17 @@ SELECT d, h, mi, s, FROM test_mdi DISTRIBUTE BY d +-- null top-level, struct, and map interval output through native shuffle +query +SELECT i, st, m +FROM VALUES + (1, + CAST(NULL AS INTERVAL DAY TO SECOND), + named_struct('i', CAST(NULL AS INTERVAL DAY TO SECOND)), + map('i', CAST(NULL AS INTERVAL DAY TO SECOND))) +AS t(k, i, st, m) +DISTRIBUTE BY k + -- overflow: days * MICROS_PER_DAY exceeds the int64 microsecond range. makeDayTimeInterval throws -- unconditionally (not ANSI-gated); this confirms the dispatched codegen path propagates Spark's -- exception. The pattern is the lowercase word so it matches every version: Spark 4.x raises diff --git a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql index b60e4efc52..e2763cce7a 100644 --- a/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql +++ b/spark/src/test/resources/sql-tests/expressions/datetime/make_ym_interval.sql @@ -53,6 +53,17 @@ SELECT y, m, FROM test_myi DISTRIBUTE BY y +-- null top-level, struct, and map interval output through native shuffle +query +SELECT i, s, m +FROM VALUES + (1, + CAST(NULL AS INTERVAL YEAR TO MONTH), + named_struct('i', CAST(NULL AS INTERVAL YEAR TO MONTH)), + map('i', CAST(NULL AS INTERVAL YEAR TO MONTH))) +AS t(k, i, s, m) +DISTRIBUTE BY k + -- overflow: years * 12 exceeds Int range. makeYearMonthInterval throws unconditionally (not -- ANSI-gated); this confirms the dispatched codegen path propagates Spark's exception. The -- pattern is the lowercase word so it matches every version: Spark 4.x raises