From a8daccfa6c79ed306cf982cd2a909dd87772f42b Mon Sep 17 00:00:00 2001 From: Rajesh Vakkalagadda Date: Fri, 18 Sep 2026 22:13:34 -0700 Subject: [PATCH 1/3] [SPARK-57812][SQL] Support catalog column-statistics serialization for nanosecond-precision timestamps CatalogColumnStat.toExternalString/fromExternalString only handled TimestampType/TimestampNTZType, so ANALYZE TABLE ... FOR COLUMNS on a TIMESTAMP_LTZ(p)/TIMESTAMP_NTZ(p) column threw columnStatisticsSerializationNotSupportedError. Add nanosecond-aware formatter cases mirroring the existing microsecond path, reusing the already-shipped TimestampFormatter.{formatNanos,parseNanos, formatWithoutTimeZoneNanos,parseWithoutTimeZoneNanos} API. Unblocking stats collection surfaced two real crashes in code that receives those stats and had never seen a nanos timestamp before: - EstimationUtils.toDouble/fromDouble had no case for the nanos types (or, pre-existing, for TimestampNTZType), so JoinEstimation, FilterEstimation.evaluateEquality/evaluateBinaryForTwoColumns, and ValueInterval all threw MatchError once CBO stats existed for such a column. FilterEstimation.evaluateBinary/evaluateInSet have their own independent type dispatch and needed the same types added directly. - CommandUtils.supportsHistogram used a broad `_: DatetimeType` match that already covered the nanos types, so ANALYZE with histograms enabled tried to run ApproximatePercentile/ ApproxCountDistinctForIntervals on a TimestampNanosVal and failed their type checks. Excluded nanos types from histogram collection instead of teaching percentile/histogram math a new composite value type; basic min/max/ndv stats are unaffected. Also fixes DESCRIBE TABLE EXTENDED showing nanosecond LTZ column stats in raw UTC instead of the session time zone, unlike its microsecond sibling. Tests: new CatalogColumnStatSuite (formatter round-trip/truncation), and new StatisticsCollectionSuite cases for the DESC round-trip, the histogram skip, and CBO estimation over nanosecond predicates (join key, equality, IN-list, two-column, and range comparisons). --- .../sql/catalyst/catalog/interface.scala | 15 +++ .../statsEstimation/EstimationUtils.scala | 12 ++- .../statsEstimation/FilterEstimation.scala | 6 +- .../catalog/CatalogColumnStatSuite.scala | 62 ++++++++++++ .../sql/execution/command/CommandUtils.scala | 5 + .../spark/sql/execution/command/tables.scala | 11 +++ .../spark/sql/StatisticsCollectionSuite.scala | 95 +++++++++++++++++++ 7 files changed, 202 insertions(+), 4 deletions(-) create mode 100644 sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/catalog/CatalogColumnStatSuite.scala diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/interface.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/interface.scala index 377b649aaca04..974b2a2ee514e 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/interface.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/interface.scala @@ -49,6 +49,7 @@ import org.apache.spark.sql.errors.QueryCompilationErrors import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ import org.apache.spark.sql.util.{CaseInsensitiveStringMap, SchemaUtils} +import org.apache.spark.unsafe.types.TimestampNanosVal import org.apache.spark.util.ArrayImplicits._ import org.apache.spark.util.Utils @@ -1047,6 +1048,13 @@ object CatalogColumnStat extends Logging { case TimestampType => getTimestampFormatter(isParsing = true).parse(s) case TimestampNTZType => getTimestampFormatter(isParsing = true, forTimestampNTZ = true).parse(s) + case t: TimestampLTZNanosType => + getTimestampFormatter(isParsing = true, format = "yyyy-MM-dd HH:mm:ss.SSSSSSSSS") + .parseNanos(s, t.precision) + case t: TimestampNTZNanosType => + getTimestampFormatter( + isParsing = true, format = "yyyy-MM-dd HH:mm:ss.SSSSSSSSS", forTimestampNTZ = true) + .parseWithoutTimeZoneNanos(s, t.precision) case ByteType => s.toByte case ShortType => s.toShort case IntegerType => s.toInt @@ -1073,6 +1081,13 @@ object CatalogColumnStat extends Logging { case TimestampNTZType => getTimestampFormatter(isParsing = false, forTimestampNTZ = true) .format(v.asInstanceOf[Long]) + case t: TimestampLTZNanosType => + getTimestampFormatter(isParsing = false, format = "yyyy-MM-dd HH:mm:ss.SSSSSSSSS") + .formatNanos(v.asInstanceOf[TimestampNanosVal], t.precision) + case t: TimestampNTZNanosType => + getTimestampFormatter( + isParsing = false, format = "yyyy-MM-dd HH:mm:ss.SSSSSSSSS", forTimestampNTZ = true) + .formatWithoutTimeZoneNanos(v.asInstanceOf[TimestampNanosVal], t.precision) case BooleanType | _: IntegralType | FloatType | DoubleType => v case _: DecimalType => v.asInstanceOf[Decimal].toJavaBigDecimal // This version of Spark does not use min/max for binary/string types so we ignore it. diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala index c107e54288529..012d8f3251611 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala @@ -23,6 +23,7 @@ import scala.math.BigDecimal.RoundingMode import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute, AttributeMap, EmptyRow, Expression} import org.apache.spark.sql.catalyst.plans.logical._ import org.apache.spark.sql.types.{DecimalType, _} +import org.apache.spark.unsafe.types.TimestampNanosVal object EstimationUtils { @@ -135,8 +136,14 @@ object EstimationUtils { */ def toDouble(value: Any, dataType: DataType): Double = { dataType match { - case _: NumericType | DateType | TimestampType => value.toString.toDouble + case _: NumericType | DateType | TimestampType | TimestampNTZType => + value.toString.toDouble case BooleanType => if (value.asInstanceOf[Boolean]) 1 else 0 + // TimestampNanosVal isn't a Long like the other datetime types, so it can't go through + // the toString/toDouble conversion above; approximate it by its epoch-microseconds + // component (dropping the sub-microsecond remainder, which is negligible for selectivity + // estimation purposes). + case _: AnyTimestampNanoType => value.asInstanceOf[TimestampNanosVal].epochMicros.toDouble } } @@ -144,7 +151,8 @@ object EstimationUtils { dataType match { case BooleanType => double.toInt == 1 case DateType => double.toInt - case TimestampType => double.toLong + case TimestampType | TimestampNTZType => double.toLong + case _: AnyTimestampNanoType => TimestampNanosVal.fromParts(double.toLong, 0.toShort) case ByteType => double.toByte case ShortType => double.toShort case IntegerType => double.toInt diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/FilterEstimation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/FilterEstimation.scala index 125ac22cbc3d8..7f986f8673bdb 100755 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/FilterEstimation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/FilterEstimation.scala @@ -291,7 +291,8 @@ case class FilterEstimation(plan: Filter) extends Logging { } attr.dataType match { - case _: NumericType | DateType | TimestampType | BooleanType => + case _: NumericType | DateType | TimestampType | TimestampNTZType | BooleanType | + _: AnyTimestampNanoType => evaluateBinaryForNumeric(op, attr, literal, update) case _: StringType | BinaryType => // TODO: It is difficult to support other binary comparisons for String/Binary @@ -413,7 +414,8 @@ case class FilterEstimation(plan: Filter) extends Logging { // use [min, max] to filter the original hSet dataType match { - case _: NumericType | BooleanType | DateType | TimestampType => + case _: NumericType | BooleanType | DateType | TimestampType | TimestampNTZType | + _: AnyTimestampNanoType => if (ndv.toDouble == 0) { return Some(0.0) } diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/catalog/CatalogColumnStatSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/catalog/CatalogColumnStatSuite.scala new file mode 100644 index 0000000000000..a7218f4b4ca36 --- /dev/null +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/catalog/CatalogColumnStatSuite.scala @@ -0,0 +1,62 @@ +/* + * 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.spark.sql.catalyst.catalog + +import org.apache.spark.SparkFunSuite +import org.apache.spark.sql.catalyst.util.TimestampNanosTestUtils.{foreachNanosPrecision, nanoOfSecTruncator, nanosVal} +import org.apache.spark.sql.types.{DataType, TimestampLTZNanosType, TimestampNTZNanosType, TimestampNTZType, TimestampType} + +class CatalogColumnStatSuite extends SparkFunSuite { + + test("SPARK-57812: nanosecond timestamp min/max round-trip through catalog stats") { + // 1970-01-01 00:00:00.123456789 on the UTC grid. + val epochMicros = 123456L + val fullNanoOfSec = 123456789 + val value = nanosVal(epochMicros, 789) + + foreachNanosPrecision { precision => + // Sub-precision digits are truncated (floored) on both format and parse, matching the + // truncation rule used by the underlying formatter. + val truncatedNanoOfSec = nanoOfSecTruncator(precision)(fullNanoOfSec) + val expected = nanosVal(epochMicros, truncatedNanoOfSec % 1000) + val expectedString = f"1970-01-01 00:00:00.$truncatedNanoOfSec%09d" + + Seq[DataType]( + TimestampLTZNanosType(precision), + TimestampNTZNanosType(precision)).foreach { dataType => + val external = CatalogColumnStat.toExternalString(value, "c", dataType) + assert(external === expectedString) + assert( + CatalogColumnStat.fromExternalString( + external, "c", dataType, CatalogColumnStat.VERSION) === expected) + } + } + } + + test("SPARK-57812: microsecond timestamp stats format is unchanged") { + // 1970-01-01 00:00:00.123456 UTC, the format ANALYZE TABLE has always persisted. + assert(CatalogColumnStat.toExternalString(123456L, "c", TimestampType) === + "1970-01-01 00:00:00.123456") + assert(CatalogColumnStat.toExternalString(123456L, "c", TimestampNTZType) === + "1970-01-01 00:00:00.123456") + assert(CatalogColumnStat.fromExternalString( + "1970-01-01 00:00:00.123456", "c", TimestampType, CatalogColumnStat.VERSION) === 123456L) + assert(CatalogColumnStat.fromExternalString( + "1970-01-01 00:00:00.123456", "c", TimestampNTZType, CatalogColumnStat.VERSION) === 123456L) + } +} diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CommandUtils.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CommandUtils.scala index 3eb2b14539ffe..3d8e513e4de54 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CommandUtils.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CommandUtils.scala @@ -371,6 +371,11 @@ object CommandUtils extends Logging { case _: IntegralType => true case _: DecimalType => true case DoubleType | FloatType => true + // ApproximatePercentile/ApproxCountDistinctForIntervals only know how to summarize types + // whose internal value is directly numeric; TimestampNanosVal (epochMicros + + // nanosWithinMicro) isn't, so nanosecond timestamps are excluded here the same way + // binary/string types are. Basic stats (min/max/ndv/nullCount) are unaffected. + case _: AnyTimestampNanoType => false case _: DatetimeType => true case _ => false } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/tables.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/tables.scala index 4d513b4d3b06f..83c7378a24e06 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/tables.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/tables.scala @@ -55,6 +55,7 @@ import org.apache.spark.sql.internal.{HiveSerDe, SQLConf} import org.apache.spark.sql.types._ import org.apache.spark.sql.util.PartitioningUtils import org.apache.spark.sql.util.SchemaUtils +import org.apache.spark.unsafe.types.TimestampNanosVal import org.apache.spark.util.ArrayImplicits._ /** @@ -915,6 +916,16 @@ case class DescribeColumnCommand( .getTimestampFormatter( isParsing = false, format = "yyyy-MM-dd HH:mm:ss.SSSSSS Z", zoneId = curZoneId) .format(internalValue.asInstanceOf[Long]) + case t: TimestampLTZNanosType => + // Same rationale as the TimestampType case above: catalog storage is always UTC, so + // convert to internal value first, then format in the current session time zone. + val internalValue = + CatalogColumnStat.fromExternalString(valueStr, name, dataType, CatalogColumnStat.VERSION) + val curZoneId = DateTimeUtils.getZoneId(SQLConf.get.sessionLocalTimeZone) + CatalogColumnStat + .getTimestampFormatter( + isParsing = false, format = "yyyy-MM-dd HH:mm:ss.SSSSSSSSS Z", zoneId = curZoneId) + .formatNanos(internalValue.asInstanceOf[TimestampNanosVal], t.precision) case _ => valueStr } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala index 45b1f7931ea17..b9578ed529207 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala @@ -642,6 +642,101 @@ class StatisticsCollectionSuite extends StatisticsCollectionTestBase with Shared } } + test("SPARK-57812: describe column stats (min, max) for nanosecond timestamp columns") { + withDefaultTimeZone(UTC) { + val table = "nanos_stats_same_time_zone" + val ltzCol = "ltz_col" + val ntzCol = "ntz_col" + withTable(table) { + sql(s"CREATE TABLE $table ($ltzCol TIMESTAMP_LTZ(9), $ntzCol TIMESTAMP_NTZ(9)) " + + "USING parquet") + sql(s"INSERT INTO $table VALUES " + + "(TIMESTAMP_LTZ'2022-01-01 00:00:01.123456789', " + + "TIMESTAMP_NTZ'2022-01-01 00:00:01.123456789'), " + + "(TIMESTAMP_LTZ'2022-01-03 00:00:02.987654321', " + + "TIMESTAMP_NTZ'2022-01-03 00:00:02.987654321')") + sql(s"ANALYZE TABLE $table COMPUTE STATISTICS FOR ALL COLUMNS") + + // LTZ catalog storage is always UTC; DESC re-renders it in the session time zone + // (UTC here, per withDefaultTimeZone above). + checkDescTimestampColStats( + tableName = table, + timestampColumn = ltzCol, + expectedMinTimestamp = "2022-01-01 00:00:01.123456789 +0000", + expectedMaxTimestamp = "2022-01-03 00:00:02.987654321 +0000") + // NTZ is zone-independent, so DESC shows the stored wall-clock value with no offset. + checkDescTimestampColStats( + tableName = table, + timestampColumn = ntzCol, + expectedMinTimestamp = "2022-01-01 00:00:01.123456789", + expectedMaxTimestamp = "2022-01-03 00:00:02.987654321") + + // Converting nanosecond-timestamp catalog stats to plan stats must not throw. + val catalogStats = getCatalogTable(table).stats.get.colStats + val ltzPlanStat = catalogStats(ltzCol).toPlanStat(ltzCol, TimestampLTZNanosType(9)) + assert(ltzPlanStat.min.isDefined && ltzPlanStat.max.isDefined) + val ntzPlanStat = catalogStats(ntzCol).toPlanStat(ntzCol, TimestampNTZNanosType(9)) + assert(ntzPlanStat.min.isDefined && ntzPlanStat.max.isDefined) + } + } + } + + test("SPARK-57812: histogram collection is skipped, not crashed, for nanosecond timestamps") { + withSQLConf(SQLConf.HISTOGRAM_ENABLED.key -> "true") { + val table = "nanos_histogram_skip" + withTable(table) { + sql(s"CREATE TABLE $table (ltz TIMESTAMP_LTZ(9)) USING parquet") + sql(s"INSERT INTO $table VALUES " + + "(TIMESTAMP_LTZ'2022-01-01 00:00:01.123456789'), " + + "(TIMESTAMP_LTZ'2022-01-03 00:00:02.987654321')") + // ApproximatePercentile/ApproxCountDistinctForIntervals don't know TimestampNanosVal, so + // this must not attempt a histogram for `ltz`; basic stats are still collected. + sql(s"ANALYZE TABLE $table COMPUTE STATISTICS FOR COLUMNS ltz") + + val colStat = getCatalogTable(table).stats.get.colStats("ltz") + assert(colStat.histogram.isEmpty) + assert(colStat.min.isDefined && colStat.max.isDefined) + } + } + } + + test("SPARK-57812: CBO estimation does not crash on nanosecond timestamp predicates") { + withSQLConf(SQLConf.CBO_ENABLED.key -> "true") { + withTable("nanos_cbo_t1", "nanos_cbo_t2") { + sql("CREATE TABLE nanos_cbo_t1(k1 TIMESTAMP_LTZ(9), k2 TIMESTAMP_LTZ(9), " + + "n TIMESTAMP_NTZ(9)) USING parquet") + sql("CREATE TABLE nanos_cbo_t2(k TIMESTAMP_LTZ(9)) USING parquet") + sql("INSERT INTO nanos_cbo_t1 VALUES " + + "(TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789', TIMESTAMP_LTZ'2022-01-02 00:00:00.1', " + + "TIMESTAMP_NTZ'2022-01-01 00:00:00.123456789'), " + + "(TIMESTAMP_LTZ'2022-01-03 00:00:00.987654321', TIMESTAMP_LTZ'2022-01-04 00:00:00.1', " + + "TIMESTAMP_NTZ'2022-01-03 00:00:00.987654321')") + sql("INSERT INTO nanos_cbo_t2 VALUES " + + "(TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789'), " + + "(TIMESTAMP_LTZ'2022-01-03 00:00:00.987654321')") + sql("ANALYZE TABLE nanos_cbo_t1 COMPUTE STATISTICS FOR COLUMNS k1, k2, n") + sql("ANALYZE TABLE nanos_cbo_t2 COMPUTE STATISTICS FOR COLUMNS k") + + // Each predicate shape below used to throw scala.MatchError once ANALYZE could persist + // nanosecond timestamp stats: the join key (JoinEstimation), and the range/equality/ + // IN-list/two-column predicates (FilterEstimation) all funnel through + // EstimationUtils.toDouble/fromDouble. .executedPlan forces the lazily-cached + // LogicalPlan.stats this predicate shape needs, the same way the join-based + // "Simple queries must be working, if CBO is turned on" test above does. + sql( + """ + |SELECT t1.k1 FROM nanos_cbo_t1 t1 + |JOIN nanos_cbo_t2 t2 ON t1.k1 = t2.k + |WHERE t1.k1 > TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789' + | AND t1.n = TIMESTAMP_NTZ'2022-01-03 00:00:00.987654321' + | AND t1.k1 IN (TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789', + | TIMESTAMP_LTZ'2022-01-03 00:00:00.987654321') + | AND t1.k2 > t1.k1 + """.stripMargin).queryExecution.executedPlan + } + } + } + private def getStatAttrNames(tableName: String): Set[String] = { val queryStats = spark.table(tableName).queryExecution.optimizedPlan.stats.attributeStats queryStats.map(_._1.name).toSet From 69559243725a7d2ec4a9378c133e6e0169231f70 Mon Sep 17 00:00:00 2001 From: Rajesh Vakkalagadda Date: Fri, 18 Sep 2026 23:22:54 -0700 Subject: [PATCH 2/3] [SPARK-57812][SQL][FOLLOWUP] Preserve nanosecond precision in CBO estimation, extend UNION support A second-pass review of the previous commit found that EstimationUtils .toDouble's nanosecond-timestamp case projected TimestampNanosVal down to epochMicros only, silently dropping the nanosWithinMicro remainder. Two distinct nanosecond values sharing an epochMicros would collapse to the same Double, corrupting CBO selectivity/min-max estimation (wrong evaluateBinaryForNumeric range checks, wrong evaluateInSet maxBy/minBy tie-breaks, wrong ValueInterval.intersect bounds) without crashing. Encode nanosWithinMicro as a fractional component instead (and decode it back with floor-based, sign-correct reconstruction in fromDouble) so distinct nanosecond values compare and round-trip correctly. Also: - UnionEstimation.isTypeSupported never gained the nanos types, so UNION ALL silently dropped min/max for them -- not a crash, but bad enough estimation input to make a downstream join look empty. PhysicalTimestampLTZNanosType/PhysicalTimestampNTZNanosType already define a full-precision Ordering[TimestampNanosVal], so this needed no lossy Double conversion, just widening the type match. - JoinEstimation.computeByHistogram bypassed EstimationUtils.toDouble with its own value.toString.toDouble, which isn't valid for TimestampNanosVal; switched it to the shared conversion. Tests: new EstimationUtilsSuite covering the toDouble/fromDouble precision fix (including pre-1970 dates), a new StatisticsCollection Suite case for the UNION fix, and extended the existing CBO test to also exercise TIMESTAMP_NTZ range/IN-list predicates (previously only equality was covered, leaving the incidental TimestampNTZType widening in evaluateBinary/evaluateInSet from the prior commit untested for those shapes). --- .../statsEstimation/EstimationUtils.scala | 21 +++++--- .../statsEstimation/JoinEstimation.scala | 8 +-- .../statsEstimation/UnionEstimation.scala | 2 +- .../EstimationUtilsSuite.scala | 54 +++++++++++++++++++ .../spark/sql/StatisticsCollectionSuite.scala | 27 ++++++++++ 5 files changed, 102 insertions(+), 10 deletions(-) create mode 100644 sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala index 012d8f3251611..9c058a24c15a1 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala @@ -138,12 +138,15 @@ object EstimationUtils { dataType match { case _: NumericType | DateType | TimestampType | TimestampNTZType => value.toString.toDouble - case BooleanType => if (value.asInstanceOf[Boolean]) 1 else 0 // TimestampNanosVal isn't a Long like the other datetime types, so it can't go through - // the toString/toDouble conversion above; approximate it by its epoch-microseconds - // component (dropping the sub-microsecond remainder, which is negligible for selectivity - // estimation purposes). - case _: AnyTimestampNanoType => value.asInstanceOf[TimestampNanosVal].epochMicros.toDouble + // the toString/toDouble conversion above. Encode the epoch-microseconds component as the + // integral part and the sub-microsecond remainder (always in [0, 999]) as a fractional + // part, so distinct nanosecond values compare and round-trip correctly through + // fromDouble instead of collapsing to the same Double. + case _: AnyTimestampNanoType => + val v = value.asInstanceOf[TimestampNanosVal] + v.epochMicros.toDouble + v.nanosWithinMicro.toDouble / 1000.0 + case BooleanType => if (value.asInstanceOf[Boolean]) 1 else 0 } } @@ -152,7 +155,13 @@ object EstimationUtils { case BooleanType => double.toInt == 1 case DateType => double.toInt case TimestampType | TimestampNTZType => double.toLong - case _: AnyTimestampNanoType => TimestampNanosVal.fromParts(double.toLong, 0.toShort) + case _: AnyTimestampNanoType => + // Inverse of the toDouble encoding above: floor (not truncate, to handle pre-1970 + // epochMicros correctly) to recover the integral microseconds, then round the + // remaining fraction back to a nanosWithinMicro in [0, 999]. + val epochMicros = math.floor(double).toLong + val nanosWithinMicro = math.round((double - epochMicros) * 1000).toShort + TimestampNanosVal.fromParts(epochMicros, nanosWithinMicro) case ByteType => double.toByte case ShortType => double.toShort case IntegerType => double.toInt diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/JoinEstimation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/JoinEstimation.scala index a4672f9cd9f6b..da3c215d12621 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/JoinEstimation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/JoinEstimation.scala @@ -263,9 +263,11 @@ case class JoinEstimation(join: Join) extends Logging { val overlappedRanges = getOverlappedRanges( leftHistogram = leftHistogram, rightHistogram = rightHistogram, - // Only numeric values have equi-height histograms. - lowerBound = newMin.get.toString.toDouble, - upperBound = newMax.get.toString.toDouble) + // Only numeric values have equi-height histograms. Go through the shared toDouble + // conversion (rather than value.toString.toDouble) so this stays correct for any type + // toDouble supports, not just ones whose internal value happens to stringify as a number. + lowerBound = toDouble(newMin.get, leftKey.dataType), + upperBound = toDouble(newMax.get, leftKey.dataType)) var card: BigDecimal = 0 var totalNdv: Double = 0 diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/UnionEstimation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/UnionEstimation.scala index 7ad05ee3ad6b3..e451a21065daf 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/UnionEstimation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/UnionEstimation.scala @@ -38,7 +38,7 @@ object UnionEstimation { private def isTypeSupported(dt: DataType): Boolean = dt match { case ByteType | IntegerType | ShortType | FloatType | LongType | DoubleType | DateType | _: DecimalType | TimestampType | TimestampNTZType | - _: AnsiIntervalType => true + _: AnsiIntervalType | _: AnyTimestampNanoType => true case _ => false } diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala new file mode 100644 index 0000000000000..c8440a0c129b6 --- /dev/null +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala @@ -0,0 +1,54 @@ +/* + * 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.spark.sql.catalyst.statsEstimation + +import org.apache.spark.SparkFunSuite +import org.apache.spark.sql.catalyst.plans.logical.statsEstimation.EstimationUtils +import org.apache.spark.sql.types.{DataType, TimestampLTZNanosType, TimestampNTZNanosType} +import org.apache.spark.unsafe.types.TimestampNanosVal + +class EstimationUtilsSuite extends SparkFunSuite { + + private val nanosTypes: Seq[DataType] = + Seq(TimestampLTZNanosType(9), TimestampNTZNanosType(9)) + + test("SPARK-57812: toDouble/fromDouble distinguish nanosecond values sharing an epochMicros") { + nanosTypes.foreach { dataType => + val low = TimestampNanosVal.fromParts(100L, 5.toShort) + val high = TimestampNanosVal.fromParts(100L, 900.toShort) + + // Both values share the same epochMicros; only an epochMicros-only projection would + // collapse them to the same Double, which previously broke ordering/tie-breaking + // (e.g. IN-list min/max, join-key interval intersection) for CBO estimation. + val lowAsDouble = EstimationUtils.toDouble(low, dataType) + val highAsDouble = EstimationUtils.toDouble(high, dataType) + assert(lowAsDouble < highAsDouble) + + assert(EstimationUtils.fromDouble(lowAsDouble, dataType) === low) + assert(EstimationUtils.fromDouble(highAsDouble, dataType) === high) + } + } + + test("SPARK-57812: toDouble/fromDouble round-trip pre-1970 nanosecond timestamps") { + nanosTypes.foreach { dataType => + val value = TimestampNanosVal.fromParts(-100L, 500.toShort) + assert( + EstimationUtils.fromDouble(EstimationUtils.toDouble(value, dataType), dataType) === value) + } + } +} diff --git a/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala index b9578ed529207..156bfa2b71a64 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/StatisticsCollectionSuite.scala @@ -723,12 +723,18 @@ class StatisticsCollectionSuite extends StatisticsCollectionTestBase with Shared // EstimationUtils.toDouble/fromDouble. .executedPlan forces the lazily-cached // LogicalPlan.stats this predicate shape needs, the same way the join-based // "Simple queries must be working, if CBO is turned on" test above does. + // The NTZ column `n` covers range/equality/IN-list too: evaluateBinary/evaluateInSet + // widen a shared type-dispatch match that (incidentally, pre-existing) never had a + // TimestampNTZType case either, and only equality on `n` would leave that gap untested. sql( """ |SELECT t1.k1 FROM nanos_cbo_t1 t1 |JOIN nanos_cbo_t2 t2 ON t1.k1 = t2.k |WHERE t1.k1 > TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789' + | AND t1.n > TIMESTAMP_NTZ'2022-01-01 00:00:00.123456789' | AND t1.n = TIMESTAMP_NTZ'2022-01-03 00:00:00.987654321' + | AND t1.n IN (TIMESTAMP_NTZ'2022-01-01 00:00:00.123456789', + | TIMESTAMP_NTZ'2022-01-03 00:00:00.987654321') | AND t1.k1 IN (TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789', | TIMESTAMP_LTZ'2022-01-03 00:00:00.987654321') | AND t1.k2 > t1.k1 @@ -737,6 +743,27 @@ class StatisticsCollectionSuite extends StatisticsCollectionTestBase with Shared } } + test("SPARK-57812: UNION propagates min/max stats for nanosecond timestamp columns") { + withSQLConf(SQLConf.CBO_ENABLED.key -> "true") { + withTable("nanos_union_t1", "nanos_union_t2") { + sql("CREATE TABLE nanos_union_t1(k TIMESTAMP_LTZ(9)) USING parquet") + sql("CREATE TABLE nanos_union_t2(k TIMESTAMP_LTZ(9)) USING parquet") + sql("INSERT INTO nanos_union_t1 VALUES (TIMESTAMP_LTZ'2022-01-01 00:00:00.123456789')") + sql("INSERT INTO nanos_union_t2 VALUES (TIMESTAMP_LTZ'2022-01-03 00:00:00.987654321')") + sql("ANALYZE TABLE nanos_union_t1 COMPUTE STATISTICS FOR COLUMNS k") + sql("ANALYZE TABLE nanos_union_t2 COMPUTE STATISTICS FOR COLUMNS k") + + // Without TimestampLTZNanosType in UnionEstimation.isTypeSupported, min/max here would + // silently come back None instead of the actual overlapping range -- not a crash, but + // bad enough estimation input to e.g. make a downstream join look empty. + val stats = sql("SELECT k FROM nanos_union_t1 UNION ALL SELECT k FROM nanos_union_t2") + .queryExecution.optimizedPlan.stats + assert(stats.attributeStats.nonEmpty) + assert(stats.attributeStats.values.forall(cs => cs.min.isDefined && cs.max.isDefined)) + } + } + } + private def getStatAttrNames(tableName: String): Set[String] = { val queryStats = spark.table(tableName).queryExecution.optimizedPlan.stats.attributeStats queryStats.map(_._1.name).toSet From 3740397a6093c21ab914d29d5eef6f4bf9603d0f Mon Sep 17 00:00:00 2001 From: Rajesh Vakkalagadda Date: Tue, 22 Sep 2026 11:39:12 -0700 Subject: [PATCH 3/3] [SPARK-57812][SQL] Make nanosecond CBO estimation honestly microsecond-resolution Review on #58946 found that the fractional sub-microsecond encoding added by the previous commit to EstimationUtils.toDouble/fromDouble only preserves nanosecond distinctions for epochMicros < 2^43 (~1970-04-12): a Double's 52-bit mantissa cannot also hold a 1/1000 fraction at realistic (e.g. 2022) magnitudes, where only 4 of 1000 sub-microsecond values remain distinguishable and a round-trip can land in the wrong microsecond entirely. The code comment, the PR description, and the previous commit's test names all claimed a lossless/full-precision round-trip the encoding never actually provided outside that epoch-adjacent window, and the added tests happened to use epochMicros values (100, -100) that stayed inside it. Revert the AnyTimestampNanoType case in toDouble/fromDouble to the epochMicros-only conversion used by the original commit on this branch, and document it honestly as microsecond resolution with the same |epochMicros| <= 2^53 exactness domain as the existing TimestampType/TimestampNTZType case. True nanosecond-resolution CBO estimation cannot be done through a Double at these magnitudes at all, and is left to SPARK-57839. Update EstimationUtilsSuite to match: round-trip now asserts truncation to the microsecond, and a new test pins the documented collision (two distinct nanosecond values sharing a microsecond compare equal) at a realistic 2022-magnitude epochMicros instead of only near the epoch. Also adds a monotonicity check across distinct epochMicros, which is the property CBO range/IN-list estimation actually relies on. Tests: catalyst/testOnly EstimationUtilsSuite FilterEstimationSuite JoinEstimationSuite UnionEstimationSuite CatalogColumnStatSuite (119/119), sql/testOnly StatisticsCollectionSuite CommandUtilsSuite (47/47). --- .../statsEstimation/EstimationUtils.scala | 26 ++++---- .../EstimationUtilsSuite.scala | 66 ++++++++++++++----- 2 files changed, 63 insertions(+), 29 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala index 9c058a24c15a1..8426c2109c897 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/statsEstimation/EstimationUtils.scala @@ -139,13 +139,16 @@ object EstimationUtils { case _: NumericType | DateType | TimestampType | TimestampNTZType => value.toString.toDouble // TimestampNanosVal isn't a Long like the other datetime types, so it can't go through - // the toString/toDouble conversion above. Encode the epoch-microseconds component as the - // integral part and the sub-microsecond remainder (always in [0, 999]) as a fractional - // part, so distinct nanosecond values compare and round-trip correctly through - // fromDouble instead of collapsing to the same Double. - case _: AnyTimestampNanoType => - val v = value.asInstanceOf[TimestampNanosVal] - v.epochMicros.toDouble + v.nanosWithinMicro.toDouble / 1000.0 + // the toString/toDouble conversion above; use its epochMicros component, giving this + // estimation path the same microsecond resolution (and the same |epochMicros| <= 2^53 + // exactness domain, ~1970 +/- 285 years) as the TimestampType/TimestampNTZType case + // above. A Double's 52-bit mantissa cannot also hold a distinguishable sub-microsecond + // fraction at realistic (i.e. any post-1970) timestamp magnitudes, so nanosWithinMicro is + // dropped here: distinct nanosecond values within the same microsecond are indistinguishable + // to CBO estimation and compare equal. This does not affect the persisted catalog + // statistic, which keeps full nanosecond precision (see CatalogColumnStat); true + // nanosecond-resolution CBO estimation is tracked separately in SPARK-57839. + case _: AnyTimestampNanoType => value.asInstanceOf[TimestampNanosVal].epochMicros.toDouble case BooleanType => if (value.asInstanceOf[Boolean]) 1 else 0 } } @@ -156,12 +159,9 @@ object EstimationUtils { case DateType => double.toInt case TimestampType | TimestampNTZType => double.toLong case _: AnyTimestampNanoType => - // Inverse of the toDouble encoding above: floor (not truncate, to handle pre-1970 - // epochMicros correctly) to recover the integral microseconds, then round the - // remaining fraction back to a nanosWithinMicro in [0, 999]. - val epochMicros = math.floor(double).toLong - val nanosWithinMicro = math.round((double - epochMicros) * 1000).toShort - TimestampNanosVal.fromParts(epochMicros, nanosWithinMicro) + // Inverse of the toDouble encoding above. nanosWithinMicro isn't recoverable (toDouble + // never encoded it), so this always reconstructs a value truncated to the microsecond. + TimestampNanosVal.fromParts(double.toLong, 0.toShort) case ByteType => double.toByte case ShortType => double.toShort case IntegerType => double.toInt diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala index c8440a0c129b6..7915cae892835 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/statsEstimation/EstimationUtilsSuite.scala @@ -27,28 +27,62 @@ class EstimationUtilsSuite extends SparkFunSuite { private val nanosTypes: Seq[DataType] = Seq(TimestampLTZNanosType(9), TimestampNTZNanosType(9)) - test("SPARK-57812: toDouble/fromDouble distinguish nanosecond values sharing an epochMicros") { + // 2022-01-01 00:00:00 UTC, in epoch microseconds: a magnitude representative of real-world + // data, unlike the epoch-adjacent values used below, where a Double happens to have more + // spare precision than this conversion actually promises. + private val realisticEpochMicros = 1640995200000000L + + test("SPARK-57812: toDouble/fromDouble round-trip nanosecond timestamps at microsecond " + + "resolution") { + nanosTypes.foreach { dataType => + val value = TimestampNanosVal.fromParts(100L, 5.toShort) + val truncated = TimestampNanosVal.fromParts(100L, 0.toShort) + val roundTripped = EstimationUtils.fromDouble(EstimationUtils.toDouble(value, dataType), + dataType) + assert(roundTripped === truncated) + } + } + + test("SPARK-57812: toDouble collapses distinct nanosecond values sharing an epochMicros, " + + "at realistic timestamp magnitudes") { + // This is the resolution the conversion actually provides -- see the comment on + // EstimationUtils.toDouble's AnyTimestampNanoType case. A Double's 52-bit mantissa cannot + // also hold a distinguishable sub-microsecond fraction once epochMicros grows past 2^43 + // (~1970-04-12), which covers every realistic (i.e. non-epoch-adjacent) timestamp, so two + // values sharing an epochMicros are indistinguishable to CBO estimation by design. + nanosTypes.foreach { dataType => + val low = TimestampNanosVal.fromParts(realisticEpochMicros, 1.toShort) + val high = TimestampNanosVal.fromParts(realisticEpochMicros, 999.toShort) + assert(EstimationUtils.toDouble(low, dataType) === EstimationUtils.toDouble(high, dataType)) + + val truncated = TimestampNanosVal.fromParts(realisticEpochMicros, 0.toShort) + val roundTripped = EstimationUtils.fromDouble(EstimationUtils.toDouble(low, dataType), + dataType) + assert(roundTripped === truncated) + } + } + + test("SPARK-57812: toDouble is monotone across distinct epochMicros at realistic magnitudes") { + // Ordering across microseconds -- rather than precision within one -- is what CBO + // range/IN-list estimation relies on, and this holds regardless of the resolution + // limitation documented above. nanosTypes.foreach { dataType => - val low = TimestampNanosVal.fromParts(100L, 5.toShort) - val high = TimestampNanosVal.fromParts(100L, 900.toShort) - - // Both values share the same epochMicros; only an epochMicros-only projection would - // collapse them to the same Double, which previously broke ordering/tie-breaking - // (e.g. IN-list min/max, join-key interval intersection) for CBO estimation. - val lowAsDouble = EstimationUtils.toDouble(low, dataType) - val highAsDouble = EstimationUtils.toDouble(high, dataType) - assert(lowAsDouble < highAsDouble) - - assert(EstimationUtils.fromDouble(lowAsDouble, dataType) === low) - assert(EstimationUtils.fromDouble(highAsDouble, dataType) === high) + val earlier = TimestampNanosVal.fromParts(realisticEpochMicros, 999.toShort) + val later = TimestampNanosVal.fromParts(realisticEpochMicros + 1, 0.toShort) + val earlierAsDouble = EstimationUtils.toDouble(earlier, dataType) + val laterAsDouble = EstimationUtils.toDouble(later, dataType) + assert(earlierAsDouble < laterAsDouble) } } - test("SPARK-57812: toDouble/fromDouble round-trip pre-1970 nanosecond timestamps") { + test("SPARK-57812: toDouble/fromDouble round-trip pre-1970 nanosecond timestamps at " + + "microsecond resolution") { nanosTypes.foreach { dataType => val value = TimestampNanosVal.fromParts(-100L, 500.toShort) - assert( - EstimationUtils.fromDouble(EstimationUtils.toDouble(value, dataType), dataType) === value) + val truncated = TimestampNanosVal.fromParts(-100L, 0.toShort) + val roundTripped = EstimationUtils.fromDouble(EstimationUtils.toDouble(value, dataType), + dataType) + assert(roundTripped === truncated) } } }