Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -135,7 +136,19 @@ 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
// TimestampNanosVal isn't a Long like the other datetime types, so it can't go through
// 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
}
}
Expand All @@ -144,7 +157,11 @@ object EstimationUtils {
dataType match {
case BooleanType => double.toInt == 1
case DateType => double.toInt
case TimestampType => double.toLong
case TimestampType | TimestampNTZType => double.toLong
case _: AnyTimestampNanoType =>
// 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
/*
* 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))

// 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 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 at " +
"microsecond resolution") {
nanosTypes.foreach { dataType =>
val value = TimestampNanosVal.fromParts(-100L, 500.toShort)
val truncated = TimestampNanosVal.fromParts(-100L, 0.toShort)
val roundTripped = EstimationUtils.fromDouble(EstimationUtils.toDouble(value, dataType),
dataType)
assert(roundTripped === truncated)
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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._

/**
Expand Down Expand Up @@ -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
}
Expand Down
Loading