From 178071116bde41b72baada35ad18f8fac4836c50 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dapeng=20Sun=28=E5=AD=99=E5=A4=A7=E9=B9=8F=29?= Date: Thu, 13 Aug 2026 04:06:23 +0800 Subject: [PATCH] [spark] Do not display an unknown partition statistic as a number Negative partition statistics mean unreported values according to PartitionStatistics#isKnown. Spark's DESCRIBE PARTITION and SHOW TABLE EXTENDED PARTITION paths displayed those sentinels directly, while SHOW also defaulted absent row counts and byte sizes to zero. Omit unknown values from DESCRIBE partition parameters, render an unknown creation time as UNKNOWN, and create its temporary display statistics only when the byte size is known. SHOW renders missing or negative row counts and byte sizes as UNKNOWN and omits negative statistic parameters. Cover negative, missing, fully reported, and partially reported metadata across the two display paths. --- .../PaimonShowTablePartitionCommand.scala | 20 +- .../util/PartitionStatisticsDisplay.scala | 54 +++ .../execution/PaimonDescribeTableExec.scala | 51 ++- .../UnreportedPartitionStatisticsTest.scala | 327 ++++++++++++++++++ 4 files changed, 433 insertions(+), 19 deletions(-) create mode 100644 paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/PartitionStatisticsDisplay.scala create mode 100644 paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/UnreportedPartitionStatisticsTest.scala diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonShowTablePartitionCommand.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonShowTablePartitionCommand.scala index 5d85c7bd5ee9..5cf5879970aa 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonShowTablePartitionCommand.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonShowTablePartitionCommand.scala @@ -21,6 +21,7 @@ package org.apache.paimon.spark.commands import org.apache.paimon.partition.PartitionStatistics import org.apache.paimon.spark.catalyst.Compatibility import org.apache.paimon.spark.leafnode.PaimonLeafRunnableCommand +import org.apache.paimon.spark.util.PartitionStatisticsDisplay import org.apache.spark.sql.{Row, SparkSession} import org.apache.spark.sql.catalyst.analysis.ResolvedPartitionSpec @@ -92,14 +93,21 @@ case class PaimonShowTablePartitionCommand( val metadata = partitionTable.loadPartitionMetadata(row) if (!metadata.isEmpty) { val metadataMap = metadata.asScala - results.put( - "Partition Parameters", - s"{${metadataMap.map { case (k, v) => s"$k=$v" }.mkString(", ")}}") + // Omit recognized Paimon statistic parameters with negative numeric values. + val reported = metadataMap.filterNot { + case (field, value) => PartitionStatisticsDisplay.isUnreported(field, value) + } + if (reported.nonEmpty) { + results.put( + "Partition Parameters", + s"{${reported.map { case (k, v) => s"$k=$v" }.mkString(", ")}}") + } - val fileSizeInBytes = - metadataMap.getOrElse(PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES, "0").toLong + // Render missing or negative row counts and byte sizes as unknown instead of zero. val recordCount = - metadataMap.getOrElse(PartitionStatistics.FIELD_RECORD_COUNT, "0").toLong + PartitionStatisticsDisplay.render(metadataMap, PartitionStatistics.FIELD_RECORD_COUNT) + val fileSizeInBytes = + PartitionStatisticsDisplay.render(metadataMap, PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES) results.put("Partition Statistics", s"$recordCount rows, $fileSizeInBytes bytes") } diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/PartitionStatisticsDisplay.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/PartitionStatisticsDisplay.scala new file mode 100644 index 000000000000..22f74c5087a4 --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/PartitionStatisticsDisplay.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.paimon.spark.util + +import org.apache.paimon.partition.PartitionStatistics + +import scala.util.Try + +/** Formatting helpers for Paimon partition statistics in Spark display commands. */ +object PartitionStatisticsDisplay { + + /** Label for a statistic or creation time that is not known. */ + val UNKNOWN: String = "UNKNOWN" + + /** The statistic fields Paimon puts into a Spark partition parameter map. */ + private val STATISTIC_FIELDS: Set[String] = Set( + PartitionStatistics.FIELD_RECORD_COUNT, + PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES, + PartitionStatistics.FIELD_FILE_COUNT, + PartitionStatistics.FIELD_LAST_FILE_CREATION_TIME + ) + + /** Returns true for a recognized Paimon statistic with a negative numeric value. */ + def isUnreported(field: String, value: String): Boolean = + STATISTIC_FIELDS.contains(field) && + asLong(value).exists(count => !PartitionStatistics.isKnown(count)) + + /** Returns the numeric value, or [[UNKNOWN]] when it is absent, nonnumeric, or negative. */ + def render(parameters: collection.Map[String, String], field: String): String = + parameters + .get(field) + .flatMap(asLong) + .filter(count => PartitionStatistics.isKnown(count)) + .map(_.toString) + .getOrElse(UNKNOWN) + + private def asLong(value: String): Option[Long] = Try(value.trim.toLong).toOption +} diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonDescribeTableExec.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonDescribeTableExec.scala index 50af431fd1fc..52f96d58f5b5 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonDescribeTableExec.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonDescribeTableExec.scala @@ -22,6 +22,7 @@ import org.apache.paimon.partition.PartitionStatistics import org.apache.paimon.spark.SparkTable import org.apache.paimon.spark.catalog.SparkBaseCatalog import org.apache.paimon.spark.leafnode.PaimonLeafV2CommandExec +import org.apache.paimon.spark.util.PartitionStatisticsDisplay import org.apache.paimon.spark.utils.CatalogUtils.{checkNamespace, toIdentifier} import org.apache.spark.sql.catalyst.InternalRow @@ -91,23 +92,44 @@ case class PaimonDescribeTableExec( } val dummyStorageFormat = CatalogStorageFormat(None, None, None, None, compressed = false, Map.empty) - val partParameters: Map[String, String] = Map( - PartitionStatistics.FIELD_FILE_COUNT -> partition.head.fileCount().toString, - PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES -> partition.head.fileSizeInBytes().toString, - PartitionStatistics.FIELD_LAST_FILE_CREATION_TIME -> partition.head - .lastFileCreationTime() - .toString, - PartitionStatistics.FIELD_RECORD_COUNT -> partition.head.recordCount().toString - ) - val partStats = - CatalogStatistics(partition.head.fileSizeInBytes(), Some(partition.head.recordCount())) - CatalogTablePartition( + val statistics = partition.head + // Include only reported values. Spark omits the "Partition Parameters" row for an empty map. + val partParameters: Map[String, String] = Seq( + PartitionStatistics.FIELD_FILE_COUNT -> statistics.fileCount(), + PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES -> statistics.fileSizeInBytes(), + PartitionStatistics.FIELD_LAST_FILE_CREATION_TIME -> statistics.lastFileCreationTime(), + PartitionStatistics.FIELD_RECORD_COUNT -> statistics.recordCount() + ).collect { + case (field, value) if PartitionStatistics.isKnown(value) => field -> value.toString + }.toMap + // CatalogTablePartition uses this temporary value only to render its "Partition Statistics" + // row. It cannot represent an unknown size, so omit the row when size is unknown and omit an + // unknown row count independently. + val partStats = if (PartitionStatistics.isKnown(statistics.fileSizeInBytes())) { + val rowCount = + if (PartitionStatistics.isKnown(statistics.recordCount())) { + Some(BigInt(statistics.recordCount())) + } else { + None + } + Some(CatalogStatistics(statistics.fileSizeInBytes(), rowCount)) + } else { + None + } + val partitionDetails = CatalogTablePartition( partitionSpec, dummyStorageFormat, partParameters, - partition.head.lastFileCreationTime(), + statistics.lastFileCreationTime(), -1, - Some(partStats)).toLinkedHashMap.foreach(s => rows += toCatalystRow(s._1, s._2, "")) + partStats).toLinkedHashMap + // Spark formats every createTime as a date; replace an unknown timestamp to avoid a 1969 date. + if (!PartitionStatistics.isKnown(statistics.lastFileCreationTime())) { + partitionDetails.put( + PaimonDescribeTableExec.CREATED_TIME_KEY, + PartitionStatisticsDisplay.UNKNOWN) + } + partitionDetails.foreach(s => rows += toCatalystRow(s._1, s._2, "")) rows += emptyRow() } @@ -124,6 +146,9 @@ case class PaimonDescribeTableExec( } object PaimonDescribeTableExec { + // Key for the partition creation-time row emitted by CatalogTablePartition.toLinkedHashMap. + val CREATED_TIME_KEY = "Created Time" + // This column metadata indicates the default value associated with a particular table column that // is in effect at any given time. Its value begins at the time of the initial CREATE/REPLACE // TABLE statement with DEFAULT column definition(s), if any. It then changes whenever an ALTER diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/UnreportedPartitionStatisticsTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/UnreportedPartitionStatisticsTest.scala new file mode 100644 index 000000000000..bd7b6b9356b4 --- /dev/null +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/UnreportedPartitionStatisticsTest.scala @@ -0,0 +1,327 @@ +/* + * 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.paimon.spark.sql + +import org.apache.paimon.catalog.{Catalog, CatalogLoader, DelegateCatalog, Identifier => PaimonIdentifier} +import org.apache.paimon.partition.{Partition, PartitionStatistics} +import org.apache.paimon.spark.{BaseTable, PaimonSparkTestBase, SparkCatalog} +import org.apache.paimon.table.{Table => PaimonTable} + +import org.apache.spark.SparkConf +import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.connector.catalog.{Identifier => SparkIdentifier, Table => SparkConnectorTable} + +import java.util.{List => JList, Map => JMap} + +import scala.collection.JavaConverters._ + +/** Verifies that Spark display commands do not show negative or missing statistics as measurements. */ +class UnreportedPartitionStatisticsTest extends PaimonSparkTestBase { + + private val partition = "20250915" + + override protected def sparkConf: SparkConf = + super.sparkConf + .set("spark.sql.catalog.paimon", classOf[UnreportedStatisticsSparkCatalog].getName) + + override protected def beforeEach(): Unit = { + UnreportedStatistics.reset() + super.beforeEach() + } + + override protected def afterEach(): Unit = { + try { + super.afterEach() + } finally { + UnreportedStatistics.reset() + } + } + + test("DESCRIBE PARTITION omits unreported statistics") { + val tableName = "describe_unreported_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val details = UnreportedStatistics.withAllUnreportedListing { + partitionBlockOf(tableName) + } + val rendered = details.mkString("\n") + + // Unknown statistics must not suppress the remaining partition details. + assert(details.exists(_._1 == "Partition Values"), rendered) + assert(!details.exists(_._1 == "Partition Parameters"), rendered) + assert(!details.exists(_._1 == "Partition Statistics"), rendered) + // An unreported creation time would otherwise read as a 1969 date. + assert(details.toMap.apply("Created Time") == "UNKNOWN", rendered) + assert(!details.exists { case (_, value) => value.contains("-1") }, rendered) + } + } + + test("DESCRIBE PARTITION preserves reported statistics") { + val tableName = "describe_reported_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val details = partitionBlockOf(tableName).toMap + val rendered = details.mkString("\n") + + assert( + details("Partition Parameters") + .contains(s"${PartitionStatistics.FIELD_RECORD_COUNT}=2"), + rendered) + assert(details("Partition Statistics").matches("""\d+ bytes, 2 rows"""), rendered) + assert(details("Created Time") != "UNKNOWN", rendered) + } + } + + test("SHOW TABLE EXTENDED PARTITION renders unreported statistics as UNKNOWN") { + val tableName = "show_unreported_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val information = UnreportedStatistics.withUnreportedMetadata { + partitionInformationOf(tableName) + } + + assert(information.contains("Partition Statistics: UNKNOWN rows, UNKNOWN bytes"), information) + assert(!information.contains("Partition Parameters"), information) + assert(!information.contains("-1"), information) + } + } + + test("SHOW TABLE EXTENDED PARTITION renders missing statistics as UNKNOWN") { + val tableName = "show_missing_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val information = UnreportedStatistics.withMissingMetadata( + PartitionStatistics.FIELD_RECORD_COUNT, + PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES) { + partitionInformationOf(tableName) + } + + assert(information.contains("Partition Statistics: UNKNOWN rows, UNKNOWN bytes"), information) + assert(information.contains(s"${PartitionStatistics.FIELD_FILE_COUNT}="), information) + assert(!information.contains(s"${PartitionStatistics.FIELD_RECORD_COUNT}="), information) + assert( + !information.contains(s"${PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES}="), + information) + } + } + + test("SHOW TABLE EXTENDED PARTITION preserves reported statistics") { + val tableName = "show_reported_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val information = partitionInformationOf(tableName) + + assert(information.contains(s"${PartitionStatistics.FIELD_RECORD_COUNT}=2"), information) + val statistics = """Partition Statistics: 2 rows, (\d+) bytes""".r + assert(statistics.findFirstMatchIn(information).exists(_.group(1).toLong > 0L), information) + } + } + + test("DESCRIBE PARTITION handles row count and size independently") { + val tableName = "describe_mixed_statistics" + withTable(tableName) { + createPartitionedTable(tableName) + insertTwoRows(tableName) + + val rowCountUnknown = + UnreportedStatistics.withUnreportedListing(PartitionStatistics.FIELD_RECORD_COUNT) { + partitionBlockOf(tableName) + } + val rowCountUnknownMap = rowCountUnknown.toMap + val rowCountUnknownRendered = rowCountUnknown.mkString("\n") + val rowCountUnknownParameters = rowCountUnknownMap("Partition Parameters") + assert( + rowCountUnknownParameters.contains(s"${PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES}="), + rowCountUnknownRendered) + assert( + !rowCountUnknownParameters.contains(s"${PartitionStatistics.FIELD_RECORD_COUNT}="), + rowCountUnknownRendered) + assert( + rowCountUnknownMap("Partition Statistics").matches("""\d+ bytes"""), + rowCountUnknownRendered) + + val sizeUnknown = + UnreportedStatistics.withUnreportedListing(PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES) { + partitionBlockOf(tableName) + } + val sizeUnknownMap = sizeUnknown.toMap + val sizeUnknownRendered = sizeUnknown.mkString("\n") + val sizeUnknownParameters = sizeUnknownMap("Partition Parameters") + assert( + sizeUnknownParameters.contains(s"${PartitionStatistics.FIELD_RECORD_COUNT}=2"), + sizeUnknownRendered) + assert( + !sizeUnknownParameters.contains(s"${PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES}="), + sizeUnknownRendered) + assert(!sizeUnknownMap.contains("Partition Statistics"), sizeUnknownRendered) + } + } + + private def createPartitionedTable(tableName: String): Unit = + sql(s"CREATE TABLE $tableName (id INT, dt STRING) PARTITIONED BY (dt)") + + private def insertTwoRows(tableName: String): Unit = + sql(s"INSERT INTO $tableName VALUES (1, '$partition'), (2, '$partition')") + + /** The "# Detailed Partition Information" rows of DESCRIBE, as (col_name, data_type) pairs. */ + private def partitionBlockOf(tableName: String): Seq[(String, String)] = { + val rows = sql(s"DESCRIBE FORMATTED $tableName PARTITION (dt = '$partition')") + .collect() + .map(row => (row.getString(0), row.getString(1))) + .toSeq + val header = rows.indexWhere(_._1 == "# Detailed Partition Information") + assert(header >= 0, rows.mkString("\n")) + // The block runs from the header to the blank row that closes it. + rows.drop(header + 1).takeWhile(_._1.nonEmpty) + } + + private def partitionInformationOf(tableName: String): String = + sql(s"SHOW TABLE EXTENDED IN $dbName0 LIKE '$tableName' PARTITION (dt = '$partition')") + .select("information") + .collect() + .head + .getString(0) +} + +/** Which statistics the display fixtures currently leave unreported. */ +private[sql] object UnreportedStatistics { + + sealed trait MetadataMode + case object AllUnknown extends MetadataMode + case class Missing(fields: Set[String]) extends MetadataMode + + /** The statistic fields [[Catalog#listPartitions]] currently leaves unreported. */ + @volatile var listingFields: Set[String] = Set.empty + + /** The partition metadata override applied by [[UnreportedStatisticsSparkCatalog]]. */ + @volatile var metadataMode: Option[MetadataMode] = None + + def reset(): Unit = { + listingFields = Set.empty + metadataMode = None + } + + def withAllUnreportedListing[T](body: => T): T = + withUnreportedListing( + PartitionStatistics.FIELD_RECORD_COUNT, + PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES, + PartitionStatistics.FIELD_FILE_COUNT, + PartitionStatistics.FIELD_LAST_FILE_CREATION_TIME + )(body) + + def withUnreportedListing[T](fields: String*)(body: => T): T = { + listingFields = fields.toSet + try body + finally listingFields = Set.empty + } + + def withUnreportedMetadata[T](body: => T): T = withMetadataMode(AllUnknown)(body) + + def withMissingMetadata[T](fields: String*)(body: => T): T = + withMetadataMode(Missing(fields.toSet))(body) + + private def withMetadataMode[T](mode: MetadataMode)(body: => T): T = { + metadataMode = Some(mode) + try body + finally metadataMode = None + } +} + +/** Catalog fixture that masks selected partition statistics. */ +private[sql] class UnreportedStatisticsSparkCatalog extends SparkCatalog { + + private lazy val unreporting: Catalog = new UnreportedStatisticsCatalog(super.paimonCatalog()) + + override def paimonCatalog(): Catalog = unreporting + + override def loadTable(ident: SparkIdentifier): SparkConnectorTable = { + super.loadTable(ident) match { + case table: BaseTable => + UnreportedStatistics.metadataMode match { + case Some(mode) => new OverriddenMetadataTable(table.table, mode) + case None => table + } + case table => table + } + } +} + +private[sql] class UnreportedStatisticsCatalog(wrapped: Catalog) extends DelegateCatalog(wrapped) { + + override def catalogLoader(): CatalogLoader = wrapped.catalogLoader() + + override def listPartitions(identifier: PaimonIdentifier): JList[Partition] = { + val partitions = super.listPartitions(identifier) + if (UnreportedStatistics.listingFields.isEmpty) { + partitions + } else { + partitions.asScala + .map( + partition => + new Partition( + partition.spec(), + listedStatistic(PartitionStatistics.FIELD_RECORD_COUNT, partition.recordCount()), + listedStatistic( + PartitionStatistics.FIELD_FILE_SIZE_IN_BYTES, + partition.fileSizeInBytes()), + listedStatistic(PartitionStatistics.FIELD_FILE_COUNT, partition.fileCount()), + listedStatistic( + PartitionStatistics.FIELD_LAST_FILE_CREATION_TIME, + partition.lastFileCreationTime()), + PartitionStatistics.UNKNOWN_TOTAL_BUCKETS, + partition.done() + )) + .asJava + } + } + + private def listedStatistic(field: String, value: Long): Long = + if (UnreportedStatistics.listingFields.contains(field)) { + PartitionStatistics.UNKNOWN + } else { + value + } +} + +/** Table wrapper that replaces or removes partition statistics for display tests. */ +private[sql] class OverriddenMetadataTable( + override val table: PaimonTable, + mode: UnreportedStatistics.MetadataMode) + extends BaseTable { + + override def loadPartitionMetadata(ident: InternalRow): JMap[String, String] = { + val reported = super.loadPartitionMetadata(ident).asScala + mode match { + case UnreportedStatistics.AllUnknown => + reported.map { case (field, _) => field -> PartitionStatistics.UNKNOWN.toString }.asJava + case UnreportedStatistics.Missing(fields) => + reported.filterNot { case (field, _) => fields.contains(field) }.asJava + } + } +}