This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 69406b51ea [metrics] Add totalFileCount and minAvgFileSize metrics to 
detect small file problem (#9449)
69406b51ea is described below

commit 69406b51eadc92907a8e4e06b077aa2246cef2c2
Author: junmuz <[email protected]>
AuthorDate: Sat Aug 29 13:55:46 2026 +0100

    [metrics] Add totalFileCount and minAvgFileSize metrics to detect small 
file problem (#9449)
---
 docs/docs/maintenance/metrics.md                   | 15 ++++++++++++
 .../java/org/apache/paimon/mergetree/Levels.java   |  4 ++++
 .../mergetree/compact/MergeTreeCompactManager.java |  1 +
 .../operation/metrics/CompactionMetrics.java       | 27 +++++++++++++++++++++
 .../operation/metrics/CompactionMetricsTest.java   | 28 ++++++++++++++++++++++
 5 files changed, 75 insertions(+)

diff --git a/docs/docs/maintenance/metrics.md b/docs/docs/maintenance/metrics.md
index 24ca72a533..e484257e70 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
             <td>Gauge</td>
             <td>The average total file size of all active (currently being 
written) buckets.</td>
         </tr>
+        <tr>
+            <td>maxTotalFileCount</td>
+            <td>Gauge</td>
+            <td>The maximum total file count of an active (currently being 
written) bucket.</td>
+        </tr>
+        <tr>
+            <td>avgTotalFileCount</td>
+            <td>Gauge</td>
+            <td>The average total file count of all active (currently being 
written) buckets.</td>
+        </tr>
+        <tr>
+            <td>minAvgFileSize</td>
+            <td>Gauge</td>
+            <td>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.</td>
+        </tr>
         <tr>
             <td>maxSortBufferUsedBytes</td>
             <td>Gauge</td>
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 39c2a45e70..574be5677d 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 class Levels {
                 + levels.stream().mapToLong(SortedRun::totalSize).sum();
     }
 
+    public long totalFileCount() {
+        return level0.size() + levels.stream().mapToInt(r -> 
r.files().size()).sum();
+    }
+
     public List<DataFileMeta> allFiles() {
         List<DataFileMeta> files = new ArrayList<>();
         List<LevelSortedRun> 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 708515bd01..b1bc38533b 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 @@ public class MergeTreeCompactManager extends 
CompactFutureManager {
         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 2a98e30d3e..ed6122b2b6 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 @@ public class CompactionMetrics {
 
         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 class CompactionMetrics {
         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 class CompactionMetrics {
 
         void reportTotalFileSize(long bytes);
 
+        void reportTotalFileCount(long count);
+
         void reportSortBufferMetrics(long usedBytes, long totalBytes);
 
         void unregister();
@@ -194,6 +215,7 @@ public class CompactionMetrics {
         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 class CompactionMetrics {
             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 df7265f07e..380eac4b91 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 class CompactionMetricsTest {
         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 class CompactionMetricsTest {
                         
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 class CompactionMetricsTest {
                     .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();

Reply via email to