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