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..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 @@ -319,11 +319,15 @@ public Plan plan() { manifestsResult.allManifests.stream() .mapToLong(f -> f.numAddedFiles() - f.numDeletedFiles()) .sum(); + // 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, snapshot == null ? 0 : snapshot.id(), manifests.size(), + manifestsResult.allManifests.size() - manifests.size(), allDataFiles - result.size(), result.size())); } 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..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 @@ -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"; @@ -46,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( @@ -59,12 +69,17 @@ 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, () -> 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()); @@ -81,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 2e8d15afb9d6..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 @@ -26,6 +26,7 @@ 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; @@ -33,11 +34,13 @@ public ScanStats( long duration, long scannedSnapshotId, long scannedManifests, + long skippedManifests, long skippedTableFiles, long resultedTableFiles) { this.duration = duration; this.scannedSnapshotId = scannedSnapshotId; this.scannedManifests = scannedManifests; + this.skippedManifests = skippedManifests; this.skippedTableFiles = skippedTableFiles; this.resultedTableFiles = resultedTableFiles; } @@ -52,6 +55,11 @@ protected long getScannedManifests() { return scannedManifests; } + @VisibleForTesting + protected long getSkippedManifests() { + return skippedManifests; + } + @VisibleForTesting protected long getSkippedTableFiles() { return skippedTableFiles; 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 7e651b838605..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 @@ -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,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)) } }