This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new bf6934a221 [spark] Improve postpone merge-on-read scan metrics (#8895)
bf6934a221 is described below
commit bf6934a221433c84c7146285342ccd7d13d5371a
Author: Zouxxyy <[email protected]>
AuthorDate: Wed Jul 29 08:10:47 2026 +0800
[spark] Improve postpone merge-on-read scan metrics (#8895)
---
.../org/apache/paimon/spark/PaimonMetrics.scala | 20 ++++
.../paimon/spark/PostponeMergeInputScan.scala | 59 ++++++++++-
.../apache/paimon/spark/PostponeMergeOnRead.scala | 20 +++-
.../spark/execution/PostponeMergeOnReadExec.scala | 89 ++++++++++++++--
.../paimon/spark/sql/PostponeBucketTableTest.scala | 115 +++++++++++++++++++--
5 files changed, 280 insertions(+), 23 deletions(-)
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonMetrics.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonMetrics.scala
index 85e639a842..ec217c0390 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonMetrics.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonMetrics.scala
@@ -31,6 +31,8 @@ object PaimonMetrics {
val SCANNED_MANIFESTS = "scannedManifests"
val SKIPPED_TABLE_FILES = "skippedTableFiles"
val RESULTED_TABLE_FILES = "resultedTableFiles"
+ val RESULTED_POSTPONE_FILES = "resultedPostponeFiles"
+ val NUM_POSTPONE_RECORDS = "numPostponeRecords"
// write metrics
val NUM_WRITERS = "numWriters"
@@ -166,6 +168,24 @@ case class PaimonResultedTableFilesTaskMetric(value: Long)
extends PaimonTaskMet
override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES
}
+case class PaimonResultedPostponeFilesMetric() extends PaimonSumMetric {
+ override def name(): String = PaimonMetrics.RESULTED_POSTPONE_FILES
+ override def description(): String = "number of resulted postpone files"
+}
+
+case class PaimonResultedPostponeFilesTaskMetric(value: Long) extends
PaimonTaskMetric {
+ override def name(): String = PaimonMetrics.RESULTED_POSTPONE_FILES
+}
+
+case class PaimonNumPostponeRecordsMetric() extends PaimonSumMetric {
+ override def name(): String = PaimonMetrics.NUM_POSTPONE_RECORDS
+ override def description(): String = "number of postpone records"
+}
+
+case class PaimonNumPostponeRecordsTaskMetric(value: Long) extends
PaimonTaskMetric {
+ override def name(): String = PaimonMetrics.NUM_POSTPONE_RECORDS
+}
+
// Write metrics
case class PaimonNumWritersMetric() extends PaimonSummaryMetric {
override def name(): String = PaimonMetrics.NUM_WRITERS
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeInputScan.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeInputScan.scala
index debb4e3450..f0d5cdb1fe 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeInputScan.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeInputScan.scala
@@ -20,9 +20,10 @@ package org.apache.paimon.spark
import org.apache.paimon.KeyValue
import org.apache.paimon.data.serializer.InternalRowSerializer
-import org.apache.paimon.reader.RecordReaderIterator
+import org.apache.paimon.reader.{RecordReader, RecordReaderIterator}
import org.apache.paimon.spark.PostponeMergeInputScan._
import org.apache.paimon.spark.PostponeMergeOnRead.MergePlan
+import org.apache.paimon.spark.util.SplitUtils
import org.apache.paimon.table.BucketMode
import org.apache.paimon.table.PostponeUtils.PostponeBucketRouter
import org.apache.paimon.table.source.{DataSplit, DeletionFile,
PostponeMergePlan, PostponeMergeReadBuilder, SplitSerializer}
@@ -31,9 +32,12 @@ import org.apache.paimon.utils.SerializationUtils
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.GenericInternalRow
+import org.apache.spark.sql.connector.metric.{CustomMetric, CustomTaskMetric}
import org.apache.spark.sql.connector.read.{Batch, InputPartition,
PartitionReader, PartitionReaderFactory, Scan}
import org.apache.spark.sql.types.{BinaryType, ByteType, IntegerType,
LongType, StructField, StructType}
+import java.util.concurrent.TimeUnit.NANOSECONDS
+
import scala.collection.JavaConverters._
/** Internal DSv2 scan which materializes real-split markers and routed
postpone records. */
@@ -44,7 +48,23 @@ private[spark] case class PostponeMergeInputScan(mergePlan:
MergePlan) extends S
override def toBatch: Batch =
PostponeMergeInputBatch(mergePlan.readBuilder, mergePlan.corePlan)
- override def description(): String = "PaimonPostponeMergeInput"
+ override def description(): String =
+ "Paimon Postpone Scan: read postpone files and route records by target
bucket"
+
+ override def supportedCustomMetrics(): Array[CustomMetric] = {
+ Array(
+ PaimonPartitionSizeMetric(),
+ PaimonReadBatchTimeMetric(),
+ PaimonResultedPostponeFilesMetric(),
+ PaimonNumPostponeRecordsMetric()
+ )
+ }
+
+ override def reportDriverMetrics(): Array[CustomTaskMetric] = {
+ val resultedPostponeFiles =
+
mergePlan.corePlan.postponeSplits().asScala.map(_.dataFiles().size().toLong).sum
+ Array(PaimonResultedPostponeFilesTaskMetric(resultedPostponeFiles))
+ }
}
private[spark] object PostponeMergeInputScan {
@@ -172,11 +192,18 @@ private[spark] object PostponeMergeInputScan {
private val keySerializer = new InternalRowSerializer(keyType)
private val valueSerializer = new InternalRowSerializer(mergeReadType)
private val partitionBytes =
SerializationUtils.serializeBinaryRow(split.partition())
+ private val timedReader =
+ new
TimedRecordReader[KeyValue](readBuilder.newRead().createPostponeReader(split))
private val records =
- new
RecordReaderIterator[KeyValue](readBuilder.newRead().createPostponeReader(split))
+ new RecordReaderIterator[KeyValue](timedReader)
private val current = new GenericInternalRow(
Array[Any](partitionBytes, 0, POSTPONE_RECORD, null, null, 0L, 0.toByte,
null))
private var nextWriterLocalOrder = 0L
+ private var numPostponeRecords = 0L
+
+ private lazy val partitionMetrics: Array[CustomTaskMetric] = {
+ Array(PaimonPartitionSizeTaskMetric(SplitUtils.splitSize(split)))
+ }
override def next(): Boolean = {
if (!records.hasNext) {
@@ -191,15 +218,41 @@ private[spark] object PostponeMergeInputScan {
current.setByte(ROW_KIND_ORDINAL, keyValue.valueKind().toByteValue)
current.update(VALUE_ORDINAL,
SerializationUtils.serializeBinaryRow(value))
nextWriterLocalOrder = Math.addExact(nextWriterLocalOrder, 1L)
+ numPostponeRecords = Math.addExact(numPostponeRecords, 1L)
true
}
}
override def get(): InternalRow = current
+ override def currentMetricsValues(): Array[CustomTaskMetric] = {
+ partitionMetrics ++ Array(
+ PaimonReadBatchTimeTaskMetric(timedReader.readBatchTimeMs),
+ PaimonNumPostponeRecordsTaskMetric(numPostponeRecords)
+ )
+ }
+
override def close(): Unit = records.close()
}
+ private[spark] class TimedRecordReader[T](delegate: RecordReader[T]) extends
RecordReader[T] {
+
+ private var readBatchTimeNs = 0L
+
+ override def readBatch(): RecordReader.RecordIterator[T] = {
+ val startTimeNs = System.nanoTime()
+ try {
+ delegate.readBatch()
+ } finally {
+ readBatchTimeNs += System.nanoTime() - startTimeNs
+ }
+ }
+
+ def readBatchTimeMs: Long = NANOSECONDS.toMillis(readBatchTimeNs)
+
+ override def close(): Unit = delegate.close()
+ }
+
private def bucketKey(split: DataSplit) = (split.partition(), split.bucket())
private def mergeRealSplits(splits: Seq[DataSplit]): DataSplit = {
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
index 501874afe2..5d7ad0eb4a 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
@@ -21,7 +21,7 @@ package org.apache.paimon.spark
import org.apache.paimon.CoreOptions
import org.apache.paimon.partition.PartitionPredicate
import org.apache.paimon.predicate.PredicateBuilder
-import org.apache.paimon.spark.PostponeMergeOnRead.MergePlan
+import org.apache.paimon.spark.PostponeMergeOnRead.{MergePlan, RealScanInfo}
import org.apache.paimon.table.{BucketMode, FileStoreTable, Table}
import org.apache.paimon.table.source.{PostponeMergePlan,
PostponeMergeReadBuilder}
@@ -75,7 +75,15 @@ final private[spark] class PostponeMergeOnRead(scan:
PaimonBaseScan) {
val postponeFiles =
corePlan.postponeSplits().asScala.iterator.map(_.dataFiles().size().toLong).sum
scan.ensureNoFullScan(postponeFiles)
- mergePlan = MergePlan(builder, corePlan,
scan.coreOptions.blobAsDescriptor())
+ val realScanInfo = RealScanInfo(
+ scan.table.fullName,
+ scan.description(),
+ scan
+ .reportDriverMetrics()
+ .map(metric => metric.name() -> metric.value())
+ .toMap)
+ mergePlan =
+ MergePlan(builder, corePlan, scan.coreOptions.blobAsDescriptor(),
realScanInfo)
}
mergePlan
}
@@ -123,5 +131,11 @@ private[spark] object PostponeMergeOnRead {
private[spark] case class MergePlan(
readBuilder: PostponeMergeReadBuilder,
corePlan: PostponeMergePlan,
- blobAsDescriptor: Boolean)
+ blobAsDescriptor: Boolean,
+ realScanInfo: RealScanInfo)
+
+ private[spark] case class RealScanInfo(
+ tableName: String,
+ description: String,
+ driverMetrics: Map[String, Long])
}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PostponeMergeOnReadExec.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PostponeMergeOnReadExec.scala
index 175709a94a..ac160233c5 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PostponeMergeOnReadExec.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PostponeMergeOnReadExec.scala
@@ -21,11 +21,13 @@ package org.apache.paimon.spark.execution
import org.apache.paimon.{CoreOptions, KeyValue}
import org.apache.paimon.data.{InternalRow => PaimonInternalRow}
import org.apache.paimon.reader.RecordReaderIterator
+import org.apache.paimon.spark.PaimonMetrics._
import org.apache.paimon.spark.PostponeMergeInputScan._
import org.apache.paimon.spark.PostponeMergeOnRead.MergePlan
import org.apache.paimon.spark.SparkUtils
import org.apache.paimon.spark.data.SparkInternalRow
import org.apache.paimon.spark.read.BinPackingSplits
+import org.apache.paimon.spark.util.SplitUtils
import org.apache.paimon.table.source.{DataSplit, PostponeMergePlan,
PostponeMergeReadBuilder, SplitSerializer}
import org.apache.paimon.types.{RowKind, RowType}
import org.apache.paimon.utils.{IteratorRecordReader, SerializationUtils}
@@ -35,7 +37,8 @@ import org.apache.spark.rdd.RDD
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute,
SortOrder, UnsafeProjection}
import org.apache.spark.sql.catalyst.plans.physical.{Distribution,
Partitioning, UnknownPartitioning}
-import org.apache.spark.sql.execution.{SparkPlan, UnaryExecNode}
+import org.apache.spark.sql.execution.{ExplainUtils, SparkPlan, SQLExecution,
UnaryExecNode}
+import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics}
import org.apache.spark.sql.internal.SQLConf
import org.apache.spark.sql.paimon.shims.SparkShimLoader
@@ -51,6 +54,38 @@ private[spark] case class PostponeMergeOnReadExec(
child: SparkPlan)
extends UnaryExecNode {
+ private val NUM_OUTPUT_ROWS = "numOutputRows"
+
+ override def nodeName: String = s"PaimonPostponeMergeScan
${mergePlan.realScanInfo.tableName}"
+
+ override def simpleString(maxFields: Int): String = {
+ s"$nodeName ${mergePlan.realScanInfo.description}"
+ }
+
+ override def verboseStringWithOperatorId(): String = {
+ s"""
+ |$formattedNodeName
+ |${ExplainUtils.generateFieldString("Input", child.output)}
+ |Scan: ${mergePlan.realScanInfo.description}
+ |""".stripMargin
+ }
+
+ override lazy val metrics: Map[String, SQLMetric] = {
+ Map(
+ NUM_OUTPUT_ROWS -> SQLMetrics.createMetric(sparkContext, "number of
output rows"),
+ NUM_SPLITS -> SQLMetrics.createMetric(sparkContext, "number of splits
read"),
+ PARTITION_SIZE -> SQLMetrics.createSizeMetric(sparkContext, "partition
size"),
+ READ_BATCH_TIME -> SQLMetrics.createTimingMetric(sparkContext, "read
batch time"),
+ PLANNING_DURATION -> SQLMetrics.createTimingMetric(sparkContext,
"planing duration"),
+ SCANNED_SNAPSHOT_ID -> SQLMetrics.createMetric(sparkContext, "scanned
snapshot id"),
+ SCANNED_MANIFESTS -> SQLMetrics.createMetric(sparkContext, "number of
scanned manifests"),
+ SKIPPED_TABLE_FILES -> SQLMetrics.createMetric(sparkContext, "number of
skipped table files"),
+ RESULTED_TABLE_FILES -> SQLMetrics.createMetric(
+ sparkContext,
+ "number of resulted table files")
+ )
+ }
+
override def requiredChildDistribution: Seq[Distribution] = {
Seq(
SparkShimLoader.shim.createClusteredDistribution(
@@ -79,10 +114,16 @@ private[spark] case class PostponeMergeOnReadExec(
}
override protected def doExecute(): RDD[InternalRow] = {
+ postDriverMetrics()
+
val readBuilder = mergePlan.readBuilder
val resultRowType = mergePlan.corePlan.resultReadType()
val blobAsDescriptor = mergePlan.blobAsDescriptor
val outputAttributes = output
+ val numOutputRows = longMetric(NUM_OUTPUT_ROWS)
+ val numSplits = longMetric(NUM_SPLITS)
+ val partitionSize = longMetric(PARTITION_SIZE)
+ val readBatchTime = longMetric(READ_BATCH_TIME)
child.execute().mapPartitions {
rows =>
@@ -91,10 +132,30 @@ private[spark] case class PostponeMergeOnReadExec(
rows,
readBuilder,
resultRowType,
- blobAsDescriptor)
- .map(row => unsafeProjection(row): InternalRow)
+ blobAsDescriptor,
+ numSplits,
+ partitionSize,
+ readBatchTime)
+ .map {
+ row =>
+ numOutputRows += 1L
+ unsafeProjection(row): InternalRow
+ }
}
}
+
+ private def postDriverMetrics(): Unit = {
+ val updatedMetrics = mergePlan.realScanInfo.driverMetrics.flatMap {
+ case (name, value) =>
+ metrics.get(name).map {
+ metric =>
+ metric.set(value)
+ metric
+ }
+ }.toSeq
+ val executionId =
sparkContext.getLocalProperty(SQLExecution.EXECUTION_ID_KEY)
+ SQLMetrics.postDriverMetricUpdates(sparkContext, executionId,
updatedMetrics)
+ }
}
private[spark] object PostponeMergeOnReadExec {
@@ -135,7 +196,10 @@ private[spark] object PostponeMergeOnReadExec {
rows: Iterator[InternalRow],
readBuilder: PostponeMergeReadBuilder,
resultRowType: RowType,
- blobAsDescriptor: Boolean)
+ blobAsDescriptor: Boolean,
+ numSplits: SQLMetric,
+ partitionSize: SQLMetric,
+ readBatchTime: SQLMetric)
extends Iterator[InternalRow]
with AutoCloseable {
@@ -144,6 +208,7 @@ private[spark] object PostponeMergeOnReadExec {
private val read = readBuilder.newRead().withIOManager(ioManager)
private val sparkRow = SparkInternalRow.create(resultRowType,
blobAsDescriptor)
private var currentReader: RecordReaderIterator[PaimonInternalRow] = _
+ private var currentTimedReader: TimedRecordReader[PaimonInternalRow] = _
private var nextRow: InternalRow = _
private var closed = false
@@ -192,6 +257,10 @@ private[spark] object PostponeMergeOnReadExec {
} else {
null
}
+ if (realSplit != null) {
+ numSplits += 1L
+ partitionSize += SplitUtils.splitSize(realSplit)
+ }
val postponeRecords = new Iterator[KeyValue] {
override def hasNext: Boolean = {
@@ -213,10 +282,11 @@ private[spark] object PostponeMergeOnReadExec {
}
}
- currentReader = new RecordReaderIterator[PaimonInternalRow](
+ currentTimedReader = new TimedRecordReader[PaimonInternalRow](
read.createBucketMergeReader(
realSplit,
new IteratorRecordReader[KeyValue](postponeRecords.asJava)))
+ currentReader = new
RecordReaderIterator[PaimonInternalRow](currentTimedReader)
if (bufferedRows.hasNext && sameBucket(bufferedRows.head, bucketKey)) {
throw new IllegalStateException(
"Unexpected postpone merge carrier kind " +
@@ -228,8 +298,13 @@ private[spark] object PostponeMergeOnReadExec {
private def closeCurrentReader(): Unit = {
if (currentReader != null) {
- currentReader.close()
- currentReader = null
+ try {
+ currentReader.close()
+ } finally {
+ readBatchTime += currentTimedReader.readBatchTimeMs
+ currentReader = null
+ currentTimedReader = null
+ }
}
}
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
index 1e9910bf7e..b25c03aec5 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
@@ -21,14 +21,16 @@ package org.apache.paimon.spark.sql
import org.apache.paimon.catalog.{Catalog, CatalogLoader, DelegateCatalog,
Identifier}
import
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX
import org.apache.paimon.fs.Path
-import org.apache.paimon.spark.{PaimonScan, PaimonSparkTestBase}
+import org.apache.paimon.spark.{PaimonScan, PaimonSparkTestBase,
PostponeMergeInputScan}
+import org.apache.paimon.spark.PaimonMetrics._
import org.apache.paimon.spark.commands.PaimonSparkWriter
+import org.apache.paimon.spark.execution.PostponeMergeOnReadExec
import org.apache.paimon.spark.procedure.SparkPostponeCompactProcedure
-import org.apache.paimon.table.{CatalogEnvironment, FileStoreTableFactory}
+import org.apache.paimon.table.{BucketMode, CatalogEnvironment,
FileStoreTableFactory}
import org.apache.spark.TaskContext
import org.apache.spark.sql.Row
-import org.apache.spark.sql.execution.datasources.v2.DataSourceV2ScanRelation
+import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec,
DataSourceV2ScanRelation}
import scala.collection.JavaConverters._
@@ -367,7 +369,7 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
withSparkSQLConf("spark.paimon.postpone.merge-on-read" -> "true") {
checkAnswer(sql("SELECT * FROM t ORDER BY k"), Seq(Row(1, "base-1"),
Row(2, "base-2")))
val plan = sql("SELECT * FROM
t").queryExecution.executedPlan.toString()
- assert(plan.contains("PostponeMergeOnRead"), plan)
+ assert(plan.contains("PaimonPostponeMergeScan test.t"), plan)
val aggregate = sql("SELECT count(*) FROM t")
checkAnswer(aggregate, Seq(Row(2L)))
@@ -443,8 +445,8 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
assert(realSplits.size() > 1, realSplits)
val plan = query.queryExecution.executedPlan.toString()
- assert(plan.contains("PostponeMergeOnRead"), plan)
- assert(plan.contains("PaimonPostponeMergeInput"), plan)
+ assert(plan.contains("PaimonPostponeMergeScan test.t"), plan)
+ assert(plan.contains("Paimon Postpone Scan"), plan)
assert(plan.contains("Exchange hashpartitioning"), plan)
assert(!plan.contains("PostponeArrivalOrder"), plan)
assert(plan.contains("Sort [__paimon_postpone_partition"), plan)
@@ -484,7 +486,99 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
"spark.paimon.scan.mode" -> "compacted-full") {
val compactedFull = sql("SELECT * FROM t ORDER BY k")
checkAnswer(compactedFull, Seq(Row(1, "base-1"), Row(2, "base-2")))
-
assert(!compactedFull.queryExecution.executedPlan.toString.contains("PostponeMergeOnRead"))
+ assert(
+ !compactedFull.queryExecution.executedPlan.toString
+ .contains("PaimonPostponeMergeScan test.t"))
+ }
+ }
+ }
+
+ test("Postpone bucket table: Spark merge on read metrics") {
+ withTable("t") {
+ sql("""
+ |CREATE TABLE t (
+ | k INT,
+ | v STRING
+ |) TBLPROPERTIES (
+ | 'primary-key' = 'k',
+ | 'bucket' = '-2',
+ | 'postpone.batch-write-fixed-bucket' = 'true',
+ | 'postpone.default-bucket-num' = '1',
+ | 'source.split.target-size' = '1 B'
+ |)
+ |""".stripMargin)
+
+ sql("INSERT INTO t VALUES (1, 'base-1')")
+ sql("INSERT INTO t VALUES (2, 'base-2')")
+ withSparkSQLConf("spark.paimon.postpone.batch-write-fixed-bucket" ->
"false") {
+ sql("INSERT INTO t VALUES (1, 'new-1'), (1, 'newest-1'), (3, 'new-3')")
+ }
+
+ val table = loadTable("t")
+ val postponeSplits = table
+ .newSnapshotReader()
+ .withBucket(BucketMode.POSTPONE_BUCKET)
+ .read()
+ .dataSplits()
+ val expectedPostponeFiles = postponeSplits.asScala
+ .map(_.dataFiles().size().toLong)
+ .sum
+ val realSplits = table
+ .newSnapshotReader()
+ .onlyReadRealBuckets()
+ .read()
+ .dataSplits()
+ val expectedRealFiles = realSplits.asScala
+ .map(_.dataFiles().size().toLong)
+ .sum
+ val expectedRealCarriers = realSplits.asScala
+ .groupBy(split => (split.partition(), split.bucket()))
+ .size
+ .toLong
+
+ withSparkSQLConf(
+ "spark.paimon.postpone.merge-on-read" -> "true",
+ "spark.sql.adaptive.enabled" -> "false") {
+ val query = sql("SELECT * FROM t WHERE k >= 1")
+ assert(
+ query.collect().sortBy(_.getInt(0)).toSeq ==
+ Seq(Row(1, "newest-1"), Row(2, "base-2"), Row(3, "new-3")))
+
+ val executedPlan = query.queryExecution.executedPlan
+ val inputScan = executedPlan.collectFirst {
+ case scan: BatchScanExec if
scan.scan.isInstanceOf[PostponeMergeInputScan] => scan
+ }.get
+ val merge = executedPlan.collectFirst { case exec:
PostponeMergeOnReadExec => exec }.get
+
+ assert(
+ inputScan.scan.description() ==
+ "Paimon Postpone Scan: read postpone files and route records by
target bucket")
+ assert(merge.nodeName == "PaimonPostponeMergeScan test.t")
+ assert(merge.simpleString(100).contains("PaimonScan"))
+ assert(merge.simpleString(100).contains("DataFilters"))
+
+ val inputMetrics = inputScan.metrics
+ assert(inputMetrics(RESULTED_POSTPONE_FILES).value ==
expectedPostponeFiles)
+ assert(inputMetrics(NUM_POSTPONE_RECORDS).value == 3L)
+ assert(!inputMetrics.contains(NUM_SPLITS))
+ assert(inputMetrics(PARTITION_SIZE).value > 0L)
+ assert(inputMetrics(READ_BATCH_TIME).value >= 0L)
+ assert(inputMetrics("numOutputRows").value == expectedRealCarriers +
3L)
+ assert(!inputMetrics.contains(RESULTED_TABLE_FILES))
+ assert(!inputMetrics.contains(PLANNING_DURATION))
+
+ val mergeMetrics = merge.metrics
+ assert(mergeMetrics(RESULTED_TABLE_FILES).value == expectedRealFiles)
+ assert(mergeMetrics(SCANNED_SNAPSHOT_ID).value == 3L)
+ assert(mergeMetrics(SCANNED_MANIFESTS).value > 0L)
+ assert(mergeMetrics(SKIPPED_TABLE_FILES).value >= 0L)
+ assert(mergeMetrics(PLANNING_DURATION).value >= 0L)
+ assert(mergeMetrics(NUM_SPLITS).value == expectedRealCarriers)
+ assert(mergeMetrics(PARTITION_SIZE).value > 0L)
+ assert(mergeMetrics(READ_BATCH_TIME).value >= 0L)
+ assert(mergeMetrics("numOutputRows").value == 3L)
+ assert(!mergeMetrics.contains(RESULTED_POSTPONE_FILES))
+ assert(!mergeMetrics.contains(NUM_POSTPONE_RECORDS))
}
}
}
@@ -535,7 +629,8 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
Row(5, "real-5"),
Row(6, "postpone-6")))
checkAnswer(sql("SELECT count(*) FROM t"), Seq(Row(5L)))
-
assert(query.queryExecution.executedPlan.toString.contains("PostponeMergeOnRead"))
+ assert(
+
query.queryExecution.executedPlan.toString.contains("PaimonPostponeMergeScan
test.t"))
}
}
}
@@ -660,8 +755,8 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
val mergePlan = sql("SELECT * FROM t WHERE pt =
'a'").queryExecution.executedPlan.toString
val ordinaryPlan =
sql("SELECT * FROM t WHERE pt =
'b'").queryExecution.executedPlan.toString
- assert(mergePlan.contains("PostponeMergeOnRead"), mergePlan)
- assert(ordinaryPlan.contains("PostponeMergeOnRead"), ordinaryPlan)
+ assert(mergePlan.contains("PaimonPostponeMergeScan test.t"), mergePlan)
+ assert(ordinaryPlan.contains("PaimonPostponeMergeScan test.t"),
ordinaryPlan)
}
}
}