)
+ 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,18 +137,21 @@ 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);
- scanMetrics.reportScan(scanStats);
+ scanMetrics.reportScan(new ScanStats(200, 1L, 20, 5, 25, 10));
+ scanMetrics.reportResultedFiles(1024, 100);
}
private void reportAgain(ScanMetrics scanMetrics) {
- ScanStats scanStats = new ScanStats(500, 2L, 22, 30, 8);
- scanMetrics.reportScan(scanStats);
+ scanMetrics.reportScan(new ScanStats(500, 2L, 22, 7, 30, 8));
+ scanMetrics.reportResultedFiles(2048, 200);
}
private ScanMetrics getScanMetrics() {
diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanResultedFilesMetricsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanResultedFilesMetricsTest.java
new file mode 100644
index 000000000000..315a1239ac3e
--- /dev/null
+++ b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanResultedFilesMetricsTest.java
@@ -0,0 +1,368 @@
+/*
+ * 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.operation.metrics;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.manifest.FileKind;
+import org.apache.paimon.manifest.ManifestEntry;
+import org.apache.paimon.metrics.Gauge;
+import org.apache.paimon.metrics.MetricGroup;
+import org.apache.paimon.metrics.MetricGroupImpl;
+import org.apache.paimon.metrics.MetricRegistry;
+import org.apache.paimon.operation.FileStoreScan;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.TableTestBase;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.source.ScanMode;
+import org.apache.paimon.table.source.StreamTableScan;
+import org.apache.paimon.table.source.snapshot.SnapshotReader;
+import org.apache.paimon.types.DataTypes;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Tests that {@link ScanMetrics#LAST_SCAN_RESULTED_TABLE_FILES_SIZE} and {@link
+ * ScanMetrics#LAST_SCAN_RESULTED_RECORD_COUNT} describe the files the consumer of a plan actually
+ * reads, across every read path of {@link SnapshotReader}.
+ *
+ * A DELTA plan carries both ADD and DELETE entries. A normal read takes only the ADD entries; a
+ * change read also reads the DELETE entries as its before files. The metrics must follow that
+ * choice: neither count the DELETE entries for an ordinary read, nor drop them for a change read.
+ */
+public class ScanResultedFilesMetricsTest extends TableTestBase {
+
+ // ------------------------------------------------------------------------------------------
+ // SnapshotReader paths
+ // ------------------------------------------------------------------------------------------
+
+ @Test
+ public void testBatchReadReportsAddEntries() throws Exception {
+ FileStoreTable table = createAppendTable("batch_read");
+ write(table, row(1, "a"), row(2, "b"));
+ write(table, row(3, "c"));
+
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ table.newSnapshotReader().withMetricRegistry(registry).read();
+
+ Expected expected = expectedOf(table, ScanMode.ALL, latestId(table), FileKind.ADD);
+ assertThat(expected.files).isGreaterThan(0);
+ assertResultedFiles(registry, expected.size, expected.records);
+ }
+
+ @Test
+ public void testNormalReadOfDeltaPlanExcludesDeleteEntries() throws Exception {
+ FileStoreTable table = createAppendTable("delta_normal_read");
+ write(table, row(1, "a"), row(2, "b"));
+ overwrite(table, row(9, "z"));
+ long overwriteId = latestId(table);
+
+ // premise: the overwrite's DELTA plan really holds entries of both kinds
+ Expected deleted = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.DELETE);
+ Expected added = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.ADD);
+ assertThat(deleted.files).isGreaterThan(0);
+ assertThat(added.files).isGreaterThan(0);
+
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ table.newSnapshotReader()
+ .withMetricRegistry(registry)
+ .withMode(ScanMode.DELTA)
+ .withSnapshot(overwriteId)
+ .read();
+
+ // only the ADD entries are read, so the DELETE side must not be counted
+ assertResultedFiles(registry, added.size, added.records);
+ // the existing file count metric keeps counting every entry of the plan; the asymmetry
+ // with size and records is deliberate, not a regression
+ assertThat(gauge(registry, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES))
+ .isEqualTo(added.files + deleted.files);
+ }
+
+ @Test
+ public void testChangeReadOfOverwriteIncludesBeforeFiles() throws Exception {
+ FileStoreTable table = createAppendTable("delta_change_read");
+ write(table, row(1, "a"), row(2, "b"));
+ overwrite(table, row(9, "z"));
+ long overwriteId = latestId(table);
+
+ Expected deleted = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.DELETE);
+ Expected added = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.ADD);
+ assertThat(deleted.files).isGreaterThan(0);
+
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ table.newSnapshotReader()
+ .withMetricRegistry(registry)
+ .withSnapshot(overwriteId)
+ .readChanges();
+
+ // the DELETE entries are the before files of the change and are read as well
+ assertResultedFiles(registry, added.size + deleted.size, added.records + deleted.records);
+ // and the change read reports strictly more than the normal read of the same plan
+ assertThat(gauge(registry, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE))
+ .isGreaterThan(added.size);
+ }
+
+ @Test
+ public void testChangeReadOfDeletionOnlyChange() throws Exception {
+ FileStoreTable table = createAppendTable("deletion_only");
+ write(table, row(1, "a"), row(2, "b"));
+ // an overwrite that writes nothing drops every existing file and adds none
+ overwrite(table);
+ long deletionId = latestId(table);
+
+ Expected deleted = expectedOf(table, ScanMode.DELTA, deletionId, FileKind.DELETE);
+ Expected added = expectedOf(table, ScanMode.DELTA, deletionId, FileKind.ADD);
+ assertThat(deleted.files).isGreaterThan(0);
+ assertThat(added.files).isEqualTo(0);
+
+ // a normal read of this plan reads nothing
+ CapturingMetricRegistry normal = new CapturingMetricRegistry();
+ table.newSnapshotReader()
+ .withMetricRegistry(normal)
+ .withMode(ScanMode.DELTA)
+ .withSnapshot(deletionId)
+ .read();
+ assertResultedFiles(normal, 0, 0);
+
+ // a change read of the same plan scans the old files to produce the deletes, and must not
+ // report zero while doing so
+ CapturingMetricRegistry change = new CapturingMetricRegistry();
+ table.newSnapshotReader().withMetricRegistry(change).withSnapshot(deletionId).readChanges();
+ assertResultedFiles(change, deleted.size, deleted.records);
+ assertThat(deleted.size).isGreaterThan(0);
+ }
+
+ @Test
+ public void testIncrementalDiffReportsBothSnapshots() throws Exception {
+ FileStoreTable table = createAppendTable("incremental_diff");
+ write(table, row(1, "a"));
+ long beforeId = latestId(table);
+ write(table, row(2, "b"), row(3, "c"));
+ long afterId = latestId(table);
+
+ Expected before = expectedOf(table, ScanMode.ALL, beforeId, FileKind.ADD);
+ Expected after = expectedOf(table, ScanMode.ALL, afterId, FileKind.ADD);
+
+ Snapshot beforeSnapshot = table.snapshotManager().snapshot(beforeId);
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ table.newSnapshotReader()
+ .withMetricRegistry(registry)
+ .withSnapshot(afterId)
+ .readIncrementalDiff(beforeSnapshot);
+
+ // the diff is computed from two independent scans, and the ADD entries of both are read
+ assertResultedFiles(registry, before.size + after.size, before.records + after.records);
+ }
+
+ @Test
+ public void testNoMetricRegistryDoesNotFail() throws Exception {
+ FileStoreTable table = createAppendTable("no_registry");
+ write(table, row(1, "a"));
+ overwrite(table, row(2, "b"));
+ long overwriteId = latestId(table);
+
+ // every path must keep working when no registry was attached
+ table.newSnapshotReader().read();
+ table.newSnapshotReader().withMode(ScanMode.DELTA).withSnapshot(overwriteId).read();
+ table.newSnapshotReader().withSnapshot(overwriteId).readChanges();
+ table.newSnapshotReader()
+ .withSnapshot(overwriteId)
+ .readIncrementalDiff(table.snapshotManager().snapshot(overwriteId - 1));
+ }
+
+ // ------------------------------------------------------------------------------------------
+ // streaming overwrite through FollowUpScanner
+ // ------------------------------------------------------------------------------------------
+
+ @Test
+ public void testStreamingOverwriteOnAppendTableReadsAddEntriesOnly() throws Exception {
+ FileStoreTable table =
+ createAppendTable(
+ "stream_append_overwrite",
+ CoreOptions.STREAMING_READ_APPEND_OVERWRITE.key(),
+ "true");
+ write(table, row(1, "a"), row(2, "b"));
+
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ StreamTableScan scan = table.newStreamScan();
+ scan.withMetricRegistry(registry);
+ scan.plan();
+
+ overwrite(table, row(9, "z"));
+ long overwriteId = latestId(table);
+ scan.plan();
+
+ // FollowUpScanner.getOverwriteChangesPlan takes the DELTA read() path for append tables,
+ // so only the ADD entries are read
+ Expected added = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.ADD);
+ Expected deleted = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.DELETE);
+ assertThat(deleted.files).isGreaterThan(0);
+ assertResultedFiles(registry, added.size, added.records);
+ }
+
+ @Test
+ public void testStreamingOverwriteOnPrimaryKeyTableReadsBeforeFiles() throws Exception {
+ FileStoreTable table =
+ createPrimaryKeyTable(
+ "stream_pk_overwrite", CoreOptions.STREAMING_READ_OVERWRITE.key(), "true");
+ write(table, row(1, "a"), row(2, "b"));
+
+ CapturingMetricRegistry registry = new CapturingMetricRegistry();
+ StreamTableScan scan = table.newStreamScan();
+ scan.withMetricRegistry(registry);
+ scan.plan();
+
+ overwrite(table, row(1, "z"));
+ long overwriteId = latestId(table);
+ scan.plan();
+
+ // FollowUpScanner.getOverwriteChangesPlan takes readChanges() for primary key tables, so
+ // the DELETE entries are read as before files and must be counted
+ Expected added = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.ADD);
+ Expected deleted = expectedOf(table, ScanMode.DELTA, overwriteId, FileKind.DELETE);
+ assertThat(deleted.files).isGreaterThan(0);
+ assertResultedFiles(registry, added.size + deleted.size, added.records + deleted.records);
+ }
+
+ // ------------------------------------------------------------------------------------------
+ // helpers
+ // ------------------------------------------------------------------------------------------
+
+ private FileStoreTable createAppendTable(String name, String... options) throws Exception {
+ return createTable(name, false, options);
+ }
+
+ private FileStoreTable createPrimaryKeyTable(String name, String... options) throws Exception {
+ return createTable(name, true, options);
+ }
+
+ private FileStoreTable createTable(String name, boolean primaryKey, String... options)
+ throws Exception {
+ Schema.Builder builder =
+ Schema.newBuilder()
+ .column("k", DataTypes.INT())
+ .column("v", DataTypes.STRING())
+ .option(CoreOptions.BUCKET.key(), "1");
+ if (primaryKey) {
+ builder.primaryKey("k");
+ } else {
+ builder.option(CoreOptions.BUCKET_KEY.key(), "k");
+ }
+ for (int i = 0; i < options.length; i += 2) {
+ builder.option(options[i], options[i + 1]);
+ }
+ Identifier identifier = identifier(name);
+ catalog.createTable(identifier, builder.build(), false);
+ return (FileStoreTable) catalog.getTable(identifier);
+ }
+
+ private static GenericRow row(int k, String v) {
+ return GenericRow.of(k, BinaryString.fromString(v));
+ }
+
+ private void overwrite(FileStoreTable table, GenericRow... rows) throws Exception {
+ BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder().withOverwrite();
+ try (BatchTableWrite write = writeBuilder.newWrite();
+ BatchTableCommit commit = writeBuilder.newCommit()) {
+ for (GenericRow r : rows) {
+ write.write(r);
+ }
+ commit.commit(write.prepareCommit());
+ }
+ }
+
+ private static long latestId(FileStoreTable table) {
+ return table.snapshotManager().latestSnapshotId();
+ }
+
+ /**
+ * Independent oracle: folds the entries of the given kind straight out of the raw store scan,
+ * without going through {@link SnapshotReader}.
+ */
+ private static Expected expectedOf(
+ FileStoreTable table, ScanMode mode, long snapshotId, FileKind kind) {
+ FileStoreScan.Plan plan =
+ table.store().newScan().withKind(mode).withSnapshot(snapshotId).plan();
+ List entries = plan.files(kind);
+ long size = 0L;
+ long records = 0L;
+ for (ManifestEntry entry : entries) {
+ size += entry.file().fileSize();
+ records += entry.file().rowCount();
+ }
+ return new Expected(entries.size(), size, records);
+ }
+
+ private static void assertResultedFiles(
+ CapturingMetricRegistry registry, long expectedSize, long expectedRecords) {
+ assertThat(gauge(registry, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE))
+ .as("resulted table files size")
+ .isEqualTo(expectedSize);
+ assertThat(gauge(registry, ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT))
+ .as("resulted record count")
+ .isEqualTo(expectedRecords);
+ }
+
+ @SuppressWarnings("unchecked")
+ private static long gauge(CapturingMetricRegistry registry, String name) {
+ MetricGroup group = registry.scanGroup;
+ assertThat(group).as("scan metric group was created").isNotNull();
+ Gauge gauge = (Gauge) group.getMetrics().get(name);
+ assertThat(gauge).as("gauge %s", name).isNotNull();
+ return gauge.getValue();
+ }
+
+ private static final class Expected {
+ private final long files;
+ private final long size;
+ private final long records;
+
+ private Expected(long files, long size, long records) {
+ this.files = files;
+ this.size = size;
+ this.records = records;
+ }
+ }
+
+ /** Keeps the scan metric group so the gauges can be read back. */
+ private static final class CapturingMetricRegistry implements MetricRegistry {
+
+ private MetricGroup scanGroup;
+
+ @Override
+ public MetricGroup createMetricGroup(String groupName, Map variables) {
+ MetricGroup group = new MetricGroupImpl(groupName, variables);
+ if (ScanMetrics.GROUP_NAME.equals(groupName)) {
+ scanGroup = group;
+ }
+ return group;
+ }
+ }
+}
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))
}
}