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();