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 1/2] [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)) } } From a076cb624e499e1bcedd01ede01c1f04226154bf Mon Sep 17 00:00:00 2001 From: Xiangyi Zhu <82511136+zhuxiangyi@users.noreply.github.com> Date: Thu, 17 Sep 2026 23:17:39 +0800 Subject: [PATCH 2/2] [core] Report resulted file size and record count after read-path selection The ADD-only fold at plan time assumed DELETE entries of a DELTA plan are never read. That holds for a normal read, but SnapshotReaderImpl.readChanges() takes those entries as IncrementalSplit.beforeFiles and IncrementalDiffSplitRead and IncrementalChangelogReadProvider read them, so a change read underreported the before side and a deletion-only change reported zero while scanning the old files. Fold the size and record count in SnapshotReaderImpl instead, once it has decided which entries it hands out: read() reports the ADD entries, including for a DELTA plan; readChanges() reports ADD and DELETE; readIncrementalDiff() reports the ADD entries of both snapshots. ScanMetrics gets a separate reportResultedFiles() for this, and ScanStats goes back to the pruning fields plus skippedManifests. The gauge names are unchanged, so the Spark and Flink bridges and the docs need no change. resultedTableFiles keeps counting every entry of the plan, as before. ScanResultedFilesMetricsTest exercises each path through the real reader and stream scan, with expected values folded independently from the raw store scan. On the previous commit its four change-read and diff cases fail and its four ADD-only cases pass. --- .../operation/AbstractFileStoreScan.java | 18 +- .../paimon/operation/metrics/ScanMetrics.java | 23 +- .../paimon/operation/metrics/ScanStats.java | 19 +- .../source/snapshot/SnapshotReaderImpl.java | 49 ++- .../operation/metrics/ScanMetricsTest.java | 8 +- .../metrics/ScanResultedFilesMetricsTest.java | 368 ++++++++++++++++++ 6 files changed, 436 insertions(+), 49 deletions(-) create mode 100644 paimon-core/src/test/java/org/apache/paimon/operation/metrics/ScanResultedFilesMetricsTest.java 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 58a2905853d8..4996792120de 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,7 +25,6 @@ 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; @@ -320,16 +319,9 @@ 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(); - } - } + // The size and record count of what will be read are not folded here: a DELTA plan + // holds both ADD and DELETE entries, and whether the DELETE entries get read depends + // on the consumer. SnapshotReaderImpl reports them once it has made that choice. scanMetrics.reportScan( new ScanStats( scanDuration, @@ -337,9 +329,7 @@ public Plan plan() { manifests.size(), manifestsResult.allManifests.size() - manifests.size(), allDataFiles - result.size(), - result.size(), - resultedTableFilesSize, - resultedRecordCount)); + result.size())); } 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 dfb1608bb8ac..5e8cb49ea9d1 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 @@ -50,6 +50,12 @@ public class ScanMetrics { private ScanStats latestScan; + // Reported separately from the scan itself: which entries of a plan are read depends on the + // consumer (a normal read takes the ADD entries, a change read takes DELETE entries as well), + // so the size and record count are folded after that choice is made, not at plan time. + private long latestResultedTableFilesSize; + private long latestResultedRecordCount; + public ScanMetrics(MetricRegistry registry, String tableName) { metricGroup = registry.createTableMetricGroup(GROUP_NAME, tableName); metricGroup.gauge( @@ -72,12 +78,8 @@ public ScanMetrics(MetricRegistry registry, String tableName) { 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(LAST_SCAN_RESULTED_TABLE_FILES_SIZE, () -> latestResultedTableFilesSize); + metricGroup.gauge(LAST_SCAN_RESULTED_RECORD_COUNT, () -> latestResultedRecordCount); metricGroup.gauge(MANIFEST_HIT_CACHE, () -> cacheMetrics.getHitObject().get()); metricGroup.gauge(MANIFEST_MISSED_CACHE, () -> cacheMetrics.getMissedObject().get()); metricGroup.gauge(DVMETA_HIT_CACHE, () -> dvMetaCacheMetrics.getHitObject().get()); @@ -94,6 +96,15 @@ public void reportScan(ScanStats scanStats) { durationHistogram.update(scanStats.getDuration()); } + /** + * Reports the total size and record count of the data files the consumer of the latest plan + * will actually read. Called after the reader has decided which entries it takes from the plan. + */ + public void reportResultedFiles(long tableFilesSize, long recordCount) { + latestResultedTableFilesSize = tableFilesSize; + latestResultedRecordCount = recordCount; + } + public CacheMetrics getCacheMetrics() { return cacheMetrics; } 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 40cd84619e50..2af722934b83 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 @@ -29,9 +29,6 @@ public class ScanStats { 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, @@ -39,17 +36,13 @@ public ScanStats( long scannedManifests, long skippedManifests, long skippedTableFiles, - long resultedTableFiles, - long resultedTableFilesSize, - long resultedRecordCount) { + long resultedTableFiles) { 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 @@ -77,16 +70,6 @@ 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/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java index cb338c857089..6f71840498c7 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java @@ -104,6 +104,7 @@ public class SnapshotReaderImpl implements SnapshotReader { private boolean hasNonPartitionFilter; private RecordComparator lazyPartitionComparator; private CacheMetrics dvMetaCacheMetrics; + @Nullable private ScanMetrics scanMetrics; public SnapshotReaderImpl( FileStoreScan scan, @@ -322,12 +323,34 @@ public SnapshotReader withBucketFilter(Filter bucketFilter) { @Override public SnapshotReader withMetricRegistry(MetricRegistry registry) { - ScanMetrics scanMetrics = new ScanMetrics(registry, tableName); + scanMetrics = new ScanMetrics(registry, tableName); dvMetaCacheMetrics = scanMetrics.getDvMetaCacheMetrics(); scan.withMetrics(scanMetrics); return this; } + /** + * Reports the size and record count of the data files this reader is about to hand out. Which + * entries of a plan are read is decided here, not in the scan: a normal read takes the ADD + * entries, a change read also reads the DELETE entries as its before files. Reporting from the + * reader keeps the metrics tied to what actually gets read. + */ + @SafeVarargs + private final void reportResultedFiles(List... readEntries) { + if (scanMetrics == null) { + return; + } + long tableFilesSize = 0L; + long recordCount = 0L; + for (List entries : readEntries) { + for (ManifestEntry entry : entries) { + tableFilesSize += entry.file().fileSize(); + recordCount += entry.file().rowCount(); + } + } + scanMetrics.reportResultedFiles(tableFilesSize, recordCount); + } + @Override public SnapshotReader withRowRanges(List sortedPushdownRowRanges) { scan.withRowRanges(sortedPushdownRowRanges); @@ -394,8 +417,10 @@ public Plan read() { FileStoreScan.Plan plan = scan.plan(); @Nullable Snapshot snapshot = plan.snapshot(); - Map>> grouped = - groupByPartFiles(plan.files(FileKind.ADD)); + // a normal read only takes the ADD entries, even from a DELTA plan + List addFiles = plan.files(FileKind.ADD); + reportResultedFiles(addFiles); + Map>> grouped = groupByPartFiles(addFiles); if (options.scanPlanSortPartition()) { Map>> sorted = new LinkedHashMap<>(); grouped.entrySet().stream() @@ -492,10 +517,15 @@ public Plan readChanges() { withMode(ScanMode.DELTA); FileStoreScan.Plan plan = scan.plan(); + List beforeEntries = plan.files(FileKind.DELETE); + List afterEntries = plan.files(FileKind.ADD); + // both sides are read: the DELETE entries are the before files of the change + reportResultedFiles(beforeEntries, afterEntries); + Map>> beforeFiles = - groupByPartFiles(plan.files(FileKind.DELETE)); + groupByPartFiles(beforeEntries); Map>> afterFiles = - groupByPartFiles(plan.files(FileKind.ADD)); + groupByPartFiles(afterEntries); LazyField beforeSnapshot = new LazyField<>(() -> snapshotManager.snapshot(plan.snapshot().id() - 1)); return toIncrementalPlan( @@ -604,10 +634,15 @@ private Plan toIncrementalPlan( public Plan readIncrementalDiff(Snapshot before) { withMode(ScanMode.ALL); FileStoreScan.Plan plan = scan.plan(); + List afterEntries = plan.files(FileKind.ADD); + List beforeEntries = scan.withSnapshot(before).plan().files(FileKind.ADD); + // two independent scans; the ADD entries of both snapshots are read to compute the diff + reportResultedFiles(beforeEntries, afterEntries); + Map>> afterFiles = - groupByPartFiles(plan.files(FileKind.ADD)); + groupByPartFiles(afterEntries); Map>> beforeFiles = - groupByPartFiles(scan.withSnapshot(before).plan().files(FileKind.ADD)); + groupByPartFiles(beforeEntries); TimeTravelUtil.checkRescaleBucketForIncrementalDiffQuery( tableSchema, before, beforeFiles, plan.snapshot(), afterFiles); return toIncrementalPlan( 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 a00f99b2d20b..19f359cc9a92 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 @@ -145,13 +145,13 @@ public void testMetricsAreUpdated() { } private void reportOnce(ScanMetrics scanMetrics) { - ScanStats scanStats = new ScanStats(200, 1L, 20, 5, 25, 10, 1024, 100); - 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, 7, 30, 8, 2048, 200); - 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; + } + } +}