From 83705bbd0f3c681cb70d64ab411aadc4e8e40c3d Mon Sep 17 00:00:00 2001 From: Xiangyi Zhu <82511136+zhuxiangyi@users.noreply.github.com> Date: Sun, 6 Sep 2026 00:52:09 +0800 Subject: [PATCH] [core][spark] Report skipped manifests, resulted file size and record count in scan metrics Scan metrics reported how much was read but not whether that amount was reasonable. lastScannedManifests only reported the count after manifest level filtering, so the pruning ratio could not be computed, and only file counts were reported, so the data volume of a scan could not be estimated. Add lastScanSkippedManifests, lastScanResultedTableFilesSize and lastScanResultedRecordCount, computed at the existing reporting site in AbstractFileStoreScan#plan from values that are already in memory, and expose them as Spark custom metrics. --- docs/docs/maintenance/metrics.md | 15 +++++++++ .../operation/AbstractFileStoreScan.java | 16 +++++++++- .../paimon/operation/metrics/ScanMetrics.java | 13 ++++++++ .../paimon/operation/metrics/ScanStats.java | 27 +++++++++++++++- .../operation/metrics/ScanMetricsTest.java | 25 +++++++++++++-- .../apache/paimon/spark/PaimonBaseScan.scala | 5 ++- .../apache/paimon/spark/PaimonMetrics.scala | 32 ++++++++++++++++++- .../spark/metric/SparkMetricRegistry.scala | 8 ++++- .../paimon/spark/sql/PaimonMetricTest.scala | 22 +++++++++---- 9 files changed, 149 insertions(+), 14 deletions(-) diff --git a/docs/docs/maintenance/metrics.md b/docs/docs/maintenance/metrics.md index e484257e7080..76a958321949 100644 --- a/docs/docs/maintenance/metrics.md +++ b/docs/docs/maintenance/metrics.md @@ -70,6 +70,11 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca Gauge Number of scanned manifest files in the last scan. + + lastScanSkippedManifests + Gauge + Number of manifest files skipped by manifest level filtering in the last scan. + lastScanSkippedTableFiles Gauge @@ -80,6 +85,16 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca Gauge Resulted table files in the last scan. + + lastScanResultedTableFilesSize + Gauge + Total size in bytes of the resulted table files to be read in the last scan. + + + lastScanResultedRecordCount + Gauge + Total number of records in the resulted table files to be read in the last scan. + diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java index a9ef5902ec9e..58a2905853d8 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java @@ -25,6 +25,7 @@ import org.apache.paimon.manifest.BucketFilter; import org.apache.paimon.manifest.FileEntry; import org.apache.paimon.manifest.FileEntry.Identifier; +import org.apache.paimon.manifest.FileKind; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.manifest.ManifestEntrySerializer; import org.apache.paimon.manifest.ManifestFile; @@ -319,13 +320,26 @@ public Plan plan() { manifestsResult.allManifests.stream() .mapToLong(f -> f.numAddedFiles() - f.numDeletedFiles()) .sum(); + // for DELTA and CHANGELOG scan modes the result contains both ADD and DELETE entries, + // only ADD entries will actually be read, so size and record count only count them + long resultedTableFilesSize = 0L; + long resultedRecordCount = 0L; + for (ManifestEntry entry : result) { + if (entry.kind() == FileKind.ADD) { + resultedTableFilesSize += entry.file().fileSize(); + resultedRecordCount += entry.file().rowCount(); + } + } scanMetrics.reportScan( new ScanStats( scanDuration, snapshot == null ? 0 : snapshot.id(), manifests.size(), + manifestsResult.allManifests.size() - manifests.size(), allDataFiles - result.size(), - result.size())); + result.size(), + resultedTableFilesSize, + resultedRecordCount)); } return new Plan() { diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanMetrics.java b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanMetrics.java index 92df327336c1..dfb1608bb8ac 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanMetrics.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanMetrics.java @@ -32,8 +32,12 @@ public class ScanMetrics { public static final String SCAN_DURATION = "scanDuration"; public static final String LAST_SCANNED_SNAPSHOT_ID = "lastScannedSnapshotId"; public static final String LAST_SCANNED_MANIFESTS = "lastScannedManifests"; + public static final String LAST_SCAN_SKIPPED_MANIFESTS = "lastScanSkippedManifests"; public static final String LAST_SCAN_SKIPPED_TABLE_FILES = "lastScanSkippedTableFiles"; public static final String LAST_SCAN_RESULTED_TABLE_FILES = "lastScanResultedTableFiles"; + public static final String LAST_SCAN_RESULTED_TABLE_FILES_SIZE = + "lastScanResultedTableFilesSize"; + public static final String LAST_SCAN_RESULTED_RECORD_COUNT = "lastScanResultedRecordCount"; public static final String MANIFEST_HIT_CACHE = "manifestHitCache"; public static final String MANIFEST_MISSED_CACHE = "manifestMissedCache"; public static final String DVMETA_HIT_CACHE = "dvMetaHitCache"; @@ -59,12 +63,21 @@ public ScanMetrics(MetricRegistry registry, String tableName) { metricGroup.gauge( LAST_SCANNED_MANIFESTS, () -> latestScan == null ? 0L : latestScan.getScannedManifests()); + metricGroup.gauge( + LAST_SCAN_SKIPPED_MANIFESTS, + () -> latestScan == null ? 0L : latestScan.getSkippedManifests()); metricGroup.gauge( LAST_SCAN_SKIPPED_TABLE_FILES, () -> latestScan == null ? 0L : latestScan.getSkippedTableFiles()); metricGroup.gauge( LAST_SCAN_RESULTED_TABLE_FILES, () -> latestScan == null ? 0L : latestScan.getResultedTableFiles()); + metricGroup.gauge( + LAST_SCAN_RESULTED_TABLE_FILES_SIZE, + () -> latestScan == null ? 0L : latestScan.getResultedTableFilesSize()); + metricGroup.gauge( + LAST_SCAN_RESULTED_RECORD_COUNT, + () -> latestScan == null ? 0L : latestScan.getResultedRecordCount()); metricGroup.gauge(MANIFEST_HIT_CACHE, () -> cacheMetrics.getHitObject().get()); metricGroup.gauge(MANIFEST_MISSED_CACHE, () -> cacheMetrics.getMissedObject().get()); metricGroup.gauge(DVMETA_HIT_CACHE, () -> dvMetaCacheMetrics.getHitObject().get()); diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanStats.java b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanStats.java index 2e8d15afb9d6..40cd84619e50 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanStats.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/ScanStats.java @@ -26,20 +26,30 @@ public class ScanStats { private final long duration; private final long scannedSnapshotId; private final long scannedManifests; + private final long skippedManifests; private final long skippedTableFiles; private final long resultedTableFiles; + // the unit is bytes + private final long resultedTableFilesSize; + private final long resultedRecordCount; public ScanStats( long duration, long scannedSnapshotId, long scannedManifests, + long skippedManifests, long skippedTableFiles, - long resultedTableFiles) { + long resultedTableFiles, + long resultedTableFilesSize, + long resultedRecordCount) { this.duration = duration; this.scannedSnapshotId = scannedSnapshotId; this.scannedManifests = scannedManifests; + this.skippedManifests = skippedManifests; this.skippedTableFiles = skippedTableFiles; this.resultedTableFiles = resultedTableFiles; + this.resultedTableFilesSize = resultedTableFilesSize; + this.resultedRecordCount = resultedRecordCount; } @VisibleForTesting @@ -52,6 +62,11 @@ protected long getScannedManifests() { return scannedManifests; } + @VisibleForTesting + protected long getSkippedManifests() { + return skippedManifests; + } + @VisibleForTesting protected long getSkippedTableFiles() { return skippedTableFiles; @@ -62,6 +77,16 @@ protected long getResultedTableFiles() { return resultedTableFiles; } + @VisibleForTesting + protected long getResultedTableFilesSize() { + return resultedTableFilesSize; + } + + @VisibleForTesting + protected long getResultedRecordCount() { + return resultedRecordCount; + } + @VisibleForTesting protected long getDuration() { return duration; diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanMetricsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanMetricsTest.java index 7e651b838605..a00f99b2d20b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanMetricsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanMetricsTest.java @@ -48,8 +48,11 @@ public void testGenericMetricsRegistration() { ScanMetrics.SCAN_DURATION, ScanMetrics.LAST_SCANNED_SNAPSHOT_ID, ScanMetrics.LAST_SCANNED_MANIFESTS, + ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS, ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES, + ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE, + ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT, ScanMetrics.MANIFEST_HIT_CACHE, ScanMetrics.MANIFEST_MISSED_CACHE, ScanMetrics.DVMETA_HIT_CACHE, @@ -71,20 +74,32 @@ public void testMetricsAreUpdated() { (Gauge) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_SNAPSHOT_ID); Gauge lastScannedManifests = (Gauge) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_MANIFESTS); + Gauge lastScanSkippedManifests = + (Gauge) registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS); Gauge lastScanSkippedTableFiles = (Gauge) registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES); Gauge lastScanResultedTableFiles = (Gauge) registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES); + Gauge lastScanResultedTableFilesSize = + (Gauge) + registeredGenericMetrics.get( + ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE); + Gauge lastScanResultedRecordCount = + (Gauge) + registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT); assertThat(lastScanDuration.getValue()).isEqualTo(0); assertThat(lastScannedSnapshotId.getValue()).isEqualTo(0); assertThat(scanDuration.getCount()).isEqualTo(0); assertThat(scanDuration.getStatistics().size()).isEqualTo(0); assertThat(lastScannedManifests.getValue()).isEqualTo(0); + assertThat(lastScanSkippedManifests.getValue()).isEqualTo(0); assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(0); assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(0); + assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(0); + assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(0); // report once reportOnce(scanMetrics); @@ -101,8 +116,11 @@ public void testMetricsAreUpdated() { assertThat(scanDuration.getStatistics().getMax()).isEqualTo(200); assertThat(scanDuration.getStatistics().getStdDev()).isEqualTo(0); assertThat(lastScannedManifests.getValue()).isEqualTo(20); + assertThat(lastScanSkippedManifests.getValue()).isEqualTo(5); assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(25); assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(10); + assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(1024); + assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(100); // report again reportAgain(scanMetrics); @@ -119,17 +137,20 @@ public void testMetricsAreUpdated() { assertThat(scanDuration.getStatistics().getMax()).isEqualTo(500); assertThat(scanDuration.getStatistics().getStdDev()).isCloseTo(212.132, offset(0.001)); assertThat(lastScannedManifests.getValue()).isEqualTo(22); + assertThat(lastScanSkippedManifests.getValue()).isEqualTo(7); assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(30); assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(8); + assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(2048); + assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(200); } private void reportOnce(ScanMetrics scanMetrics) { - ScanStats scanStats = new ScanStats(200, 1L, 20, 25, 10); + ScanStats scanStats = new ScanStats(200, 1L, 20, 5, 25, 10, 1024, 100); scanMetrics.reportScan(scanStats); } private void reportAgain(ScanMetrics scanMetrics) { - ScanStats scanStats = new ScanStats(500, 2L, 22, 30, 8); + ScanStats scanStats = new ScanStats(500, 2L, 22, 7, 30, 8, 2048, 200); scanMetrics.reportScan(scanStats); } diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala index dc71a9cfbfbf..9951030b5090 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala @@ -178,7 +178,10 @@ abstract class PaimonBaseScan(table: InnerTable) PaimonPlanningDurationMetric(), PaimonScannedSnapshotIdMetric(), PaimonScannedManifestsMetric(), - PaimonSkippedTableFilesMetric() + PaimonSkippedManifestsMetric(), + PaimonSkippedTableFilesMetric(), + PaimonResultedTableFilesSizeMetric(), + PaimonResultedRecordCountMetric() ) } 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 ec217c0390ad..1ea9b9860147 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 @@ -29,8 +29,11 @@ object PaimonMetrics { val PLANNING_DURATION = "planningDuration" val SCANNED_SNAPSHOT_ID = "scannedSnapshotId" val SCANNED_MANIFESTS = "scannedManifests" + val SKIPPED_MANIFESTS = "skippedManifests" val SKIPPED_TABLE_FILES = "skippedTableFiles" val RESULTED_TABLE_FILES = "resultedTableFiles" + val RESULTED_TABLE_FILES_SIZE = "resultedTableFilesSize" + val RESULTED_RECORD_COUNT = "resultedRecordCount" val RESULTED_POSTPONE_FILES = "resultedPostponeFiles" val NUM_POSTPONE_RECORDS = "numPostponeRecords" @@ -125,7 +128,7 @@ case class PaimonReadBatchTimeTaskMetric(value: Long) extends PaimonTaskMetric { case class PaimonPlanningDurationMetric() extends PaimonTimingSumMetric { override def name(): String = PaimonMetrics.PLANNING_DURATION - override def description(): String = "planing duration" + override def description(): String = "planning duration" } case class PaimonPlanningDurationTaskMetric(value: Long) extends PaimonTaskMetric { @@ -150,6 +153,15 @@ case class PaimonScannedManifestsTaskMetric(value: Long) extends PaimonTaskMetri override def name(): String = PaimonMetrics.SCANNED_MANIFESTS } +case class PaimonSkippedManifestsMetric() extends PaimonSumMetric { + override def name(): String = PaimonMetrics.SKIPPED_MANIFESTS + override def description(): String = "number of skipped manifests" +} + +case class PaimonSkippedManifestsTaskMetric(value: Long) extends PaimonTaskMetric { + override def name(): String = PaimonMetrics.SKIPPED_MANIFESTS +} + case class PaimonSkippedTableFilesMetric() extends PaimonSumMetric { override def name(): String = PaimonMetrics.SKIPPED_TABLE_FILES override def description(): String = "number of skipped table files" @@ -168,6 +180,24 @@ case class PaimonResultedTableFilesTaskMetric(value: Long) extends PaimonTaskMet override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES } +case class PaimonResultedTableFilesSizeMetric() extends PaimonSizeSumMetric { + override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES_SIZE + override def description(): String = "size of resulted table files" +} + +case class PaimonResultedTableFilesSizeTaskMetric(value: Long) extends PaimonTaskMetric { + override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES_SIZE +} + +case class PaimonResultedRecordCountMetric() extends PaimonSumMetric { + override def name(): String = PaimonMetrics.RESULTED_RECORD_COUNT + override def description(): String = "number of resulted records" +} + +case class PaimonResultedRecordCountTaskMetric(value: Long) extends PaimonTaskMetric { + override def name(): String = PaimonMetrics.RESULTED_RECORD_COUNT +} + case class PaimonResultedPostponeFilesMetric() extends PaimonSumMetric { override def name(): String = PaimonMetrics.RESULTED_POSTPONE_FILES override def description(): String = "number of resulted postpone files" diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/metric/SparkMetricRegistry.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/metric/SparkMetricRegistry.scala index 9aeeed7a03cb..863411a899b7 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/metric/SparkMetricRegistry.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/metric/SparkMetricRegistry.scala @@ -50,10 +50,16 @@ case class SparkMetricRegistry() extends MetricRegistry { gauge[Long](metrics, ScanMetrics.LAST_SCANNED_SNAPSHOT_ID)), PaimonScannedManifestsTaskMetric( gauge[Long](metrics, ScanMetrics.LAST_SCANNED_MANIFESTS)), + PaimonSkippedManifestsTaskMetric( + gauge[Long](metrics, ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS)), PaimonSkippedTableFilesTaskMetric( gauge[Long](metrics, ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES)), PaimonResultedTableFilesTaskMetric( - gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES)) + gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES)), + PaimonResultedTableFilesSizeTaskMetric( + gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE)), + PaimonResultedRecordCountTaskMetric( + gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT)) ) case None => Array.empty diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PaimonMetricTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PaimonMetricTest.scala index 3b6784189255..366e9988da97 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PaimonMetricTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PaimonMetricTest.scala @@ -18,7 +18,7 @@ package org.apache.paimon.spark.sql -import org.apache.paimon.spark.PaimonMetrics.{RESULTED_TABLE_FILES, SCANNED_SNAPSHOT_ID, SKIPPED_TABLE_FILES} +import org.apache.paimon.spark.PaimonMetrics.{RESULTED_RECORD_COUNT, RESULTED_TABLE_FILES, RESULTED_TABLE_FILES_SIZE, SCANNED_SNAPSHOT_ID, SKIPPED_MANIFESTS, SKIPPED_TABLE_FILES} import org.apache.paimon.spark.PaimonSparkTestBase import org.apache.paimon.spark.read.PaimonSplitScan import org.apache.paimon.spark.util.ScanPlanHelper @@ -53,7 +53,8 @@ class PaimonMetricTest extends PaimonSparkTestBase with ScanPlanHelper { s: String, scannedSnapshotId: Long, skippedTableFiles: Long, - resultedTableFiles: Long): Unit = { + resultedTableFiles: Long, + resultedRecordCount: Long): Unit = { val scan = getPaimonScan(s) // call getInputPartitions to trigger scan scan.inputPartitions @@ -61,17 +62,24 @@ class PaimonMetricTest extends PaimonSparkTestBase with ScanPlanHelper { Assertions.assertEquals(scannedSnapshotId, metric(metrics, SCANNED_SNAPSHOT_ID)) Assertions.assertEquals(skippedTableFiles, metric(metrics, SKIPPED_TABLE_FILES)) Assertions.assertEquals(resultedTableFiles, metric(metrics, RESULTED_TABLE_FILES)) + Assertions.assertEquals(resultedRecordCount, metric(metrics, RESULTED_RECORD_COUNT)) + Assertions.assertTrue(metric(metrics, RESULTED_TABLE_FILES_SIZE) > 0) } - checkMetrics(s"SELECT * FROM T", 3, 0, 5) - checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 3, 2, 3) + checkMetrics(s"SELECT * FROM T", 3, 0, 5, 5) + checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 3, 2, 3, 3) sql(s"DELETE FROM T WHERE pt = 'p1'") - checkMetrics(s"SELECT * FROM T", 4, 0, 4) + checkMetrics(s"SELECT * FROM T", 4, 0, 4, 4) sql("CALL sys.compact(table => 'T', partitions => 'pt=\"p2\"')") - checkMetrics(s"SELECT * FROM T", 5, 0, 2) - checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 5, 1, 1) + checkMetrics(s"SELECT * FROM T", 5, 0, 2, 4) + checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 5, 1, 1, 3) + + // a scan without any filter cannot prune any manifest + val fullScan = getPaimonScan(s"SELECT * FROM T") + fullScan.inputPartitions + Assertions.assertEquals(0, metric(fullScan.reportDriverMetrics(), SKIPPED_MANIFESTS)) } }