From 99b28b6ee85fe1f19ea936f520f3341e69dc03c7 Mon Sep 17 00:00:00 2001 From: Muhammad Junaid Muzammil <4795269+junmuz@users.noreply.github.com> Date: Thu, 27 Aug 2026 04:32:20 -0700 Subject: [PATCH] [metrics] Add totalFileCount and minAvgFileSize metrics to detect small file problem --- docs/docs/maintenance/metrics.md | 15 ++++++++++ .../org/apache/paimon/mergetree/Levels.java | 4 +++ .../compact/MergeTreeCompactManager.java | 1 + .../operation/metrics/CompactionMetrics.java | 27 ++++++++++++++++++ .../metrics/CompactionMetricsTest.java | 28 +++++++++++++++++++ 5 files changed, 75 insertions(+) diff --git a/docs/docs/maintenance/metrics.md b/docs/docs/maintenance/metrics.md index 24ca72a53369..e484257e7080 100644 --- a/docs/docs/maintenance/metrics.md +++ b/docs/docs/maintenance/metrics.md @@ -407,6 +407,21 @@ Lookup metrics are available for local partial lookup. They are reported at look Gauge The average total file size of all active (currently being written) buckets. + + maxTotalFileCount + Gauge + The maximum total file count of an active (currently being written) bucket. + + + avgTotalFileCount + Gauge + The average total file count of all active (currently being written) buckets. + + + minAvgFileSize + Gauge + The minimum average file size across all active buckets, computed as total file size divided by total file count per bucket. Directly indicates if any bucket has a small file problem. Only reported for primary-key tables. + maxSortBufferUsedBytes Gauge diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/Levels.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/Levels.java index 39c2a45e702c..574be5677db2 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/Levels.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/Levels.java @@ -150,6 +150,10 @@ public long totalFileSize() { + levels.stream().mapToLong(SortedRun::totalSize).sum(); } + public long totalFileCount() { + return level0.size() + levels.stream().mapToInt(r -> r.files().size()).sum(); + } + public List allFiles() { List files = new ArrayList<>(); List runs = levelSortedRuns(); diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java index 708515bd014f..b1bc38533bcb 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java @@ -298,6 +298,7 @@ private void reportMetrics() { if (metricsReporter != null) { metricsReporter.reportLevel0FileCount(levels.level0().size()); metricsReporter.reportTotalFileSize(levels.totalFileSize()); + metricsReporter.reportTotalFileCount(levels.totalFileCount()); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java index 2a98e30d3ec0..ed6122b2b6dc 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java @@ -51,6 +51,9 @@ public class CompactionMetrics { public static final String AVG_COMPACTION_OUTPUT_SIZE = "avgCompactionOutputSize"; public static final String MAX_TOTAL_FILE_SIZE = "maxTotalFileSize"; public static final String AVG_TOTAL_FILE_SIZE = "avgTotalFileSize"; + public static final String MAX_TOTAL_FILE_COUNT = "maxTotalFileCount"; + public static final String AVG_TOTAL_FILE_COUNT = "avgTotalFileCount"; + public static final String MIN_AVG_FILE_SIZE = "minAvgFileSize"; public static final String MAX_SORT_BUFFER_USED_BYTES = "maxSortBufferUsedBytes"; public static final String AVG_SORT_BUFFER_USED_BYTES = "avgSortBufferUsedBytes"; @@ -107,6 +110,10 @@ private void registerGenericCompactionMetrics() { metricGroup.gauge(MAX_TOTAL_FILE_SIZE, () -> getTotalFileSizeStream().max().orElse(-1)); metricGroup.gauge(AVG_TOTAL_FILE_SIZE, () -> getTotalFileSizeStream().average().orElse(-1)); + metricGroup.gauge(MAX_TOTAL_FILE_COUNT, () -> getTotalFileCountStream().max().orElse(-1)); + metricGroup.gauge( + AVG_TOTAL_FILE_COUNT, () -> getTotalFileCountStream().average().orElse(-1)); + metricGroup.gauge(MIN_AVG_FILE_SIZE, () -> getAvgFileSizeStream().min().orElse(-1)); metricGroup.gauge( MAX_SORT_BUFFER_USED_BYTES, () -> getSortBufferUsedBytesStream().max().orElse(-1)); @@ -147,6 +154,18 @@ public LongStream getTotalFileSizeStream() { return reporters.values().stream().mapToLong(r -> r.totalFileSize); } + @VisibleForTesting + public LongStream getTotalFileCountStream() { + return reporters.values().stream().mapToLong(r -> r.totalFileCount); + } + + @VisibleForTesting + public LongStream getAvgFileSizeStream() { + return reporters.values().stream() + .filter(r -> r.totalFileCount > 0) + .mapToLong(r -> r.totalFileSize / r.totalFileCount); + } + private LongStream getSortBufferUsedBytesStream() { return reporters.values().stream().mapToLong(r -> r.sortBufferUsedBytes); } @@ -182,6 +201,8 @@ public interface Reporter { void reportTotalFileSize(long bytes); + void reportTotalFileCount(long count); + void reportSortBufferMetrics(long usedBytes, long totalBytes); void unregister(); @@ -194,6 +215,7 @@ private class ReporterImpl implements Reporter { private long compactionInputSize = 0; private long compactionOutputSize = 0; private long totalFileSize = 0; + private long totalFileCount = 0; private long sortBufferUsedBytes = 0; private double sortBufferUtilisationPercent = 0.0; @@ -234,6 +256,11 @@ public void reportTotalFileSize(long bytes) { this.totalFileSize = bytes; } + @Override + public void reportTotalFileCount(long count) { + this.totalFileCount = count; + } + @Override public void reportLevel0FileCount(long count) { this.level0FileCount = count; diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java index df7265f07e80..380eac4b918a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java @@ -132,6 +132,23 @@ public void testReportMetrics() { assertThat(getMetric(metrics, CompactionMetrics.MAX_LEVEL0_FILE_COUNT)).isEqualTo(8L); assertThat(getMetric(metrics, CompactionMetrics.AVG_LEVEL0_FILE_COUNT)).isEqualTo(5.0); + reporters[0].reportTotalFileCount(10); + reporters[1].reportTotalFileCount(6); + reporters[2].reportTotalFileCount(8); + assertThat(getMetric(metrics, CompactionMetrics.MAX_TOTAL_FILE_COUNT)).isEqualTo(10L); + assertThat(getMetric(metrics, CompactionMetrics.AVG_TOTAL_FILE_COUNT)).isEqualTo(8.0); + + reporters[0].reportTotalFileCount(15); + assertThat(getMetric(metrics, CompactionMetrics.MAX_TOTAL_FILE_COUNT)).isEqualTo(15L); + assertThat(getMetric(metrics, CompactionMetrics.AVG_TOTAL_FILE_COUNT)) + .isEqualTo(29.0 / 3.0); + + // report file sizes to test minAvgFileSize + reporters[0].reportTotalFileSize(150_000_000); // 150MB / 15 files = 10MB avg + reporters[1].reportTotalFileSize(6_000_000); // 6MB / 6 files = 1MB avg (smallest) + reporters[2].reportTotalFileSize(80_000_000); // 80MB / 8 files = 10MB avg + assertThat(getMetric(metrics, CompactionMetrics.MIN_AVG_FILE_SIZE)).isEqualTo(1_000_000L); + reporters[0].reportCompactionTime(300000); reporters[0].reportCompactionTime(250000); reporters[0].reportCompactionTime(270000); @@ -216,6 +233,12 @@ public void testTotalFileSizeForPrimaryKeyTables() throws Exception { dataSplit.dataFiles().stream().mapToLong(DataFileMeta::fileSize).sum(); } + long[] totalFileCounts = new long[bucketNum]; + for (Split split : table.newScan().plan().splits()) { + DataSplit dataSplit = (DataSplit) split; + totalFileCounts[dataSplit.bucket()] += dataSplit.dataFiles().size(); + } + CompactionMetrics metrics = ((AbstractFileStoreWrite) write.getWrite()).compactionMetrics(); assertThat(metrics.getTotalFileSizeStream()).hasSize(bucketNum); @@ -223,6 +246,11 @@ public void testTotalFileSizeForPrimaryKeyTables() throws Exception { .isEqualTo(Arrays.stream(totalFileSizes).max().orElse(0)); assertThat(getMetric(metrics, CompactionMetrics.AVG_TOTAL_FILE_SIZE)) .isEqualTo(Arrays.stream(totalFileSizes).average().orElse(0)); + assertThat(metrics.getTotalFileCountStream()).hasSize(bucketNum); + assertThat(getMetric(metrics, CompactionMetrics.MAX_TOTAL_FILE_COUNT)) + .isEqualTo(Arrays.stream(totalFileCounts).max().orElse(0)); + assertThat(getMetric(metrics, CompactionMetrics.AVG_TOTAL_FILE_COUNT)) + .isEqualTo(Arrays.stream(totalFileCounts).average().orElse(0)); } write.close();