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