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)) } }