Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions docs/docs/maintenance/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -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>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<DataFileMeta> allFiles() {
List<DataFileMeta> files = new ArrayList<>();
List<LevelSortedRun> runs = levelSortedRuns();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,7 @@ private void reportMetrics() {
if (metricsReporter != null) {
metricsReporter.reportLevel0FileCount(levels.level0().size());
metricsReporter.reportTotalFileSize(levels.totalFileSize());
metricsReporter.reportTotalFileCount(levels.totalFileCount());
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -182,6 +201,8 @@ public interface Reporter {

void reportTotalFileSize(long bytes);

void reportTotalFileCount(long count);

void reportSortBufferMetrics(long usedBytes, long totalBytes);

void unregister();
Expand All @@ -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;

Expand Down Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -216,13 +233,24 @@ 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);
assertThat(getMetric(metrics, CompactionMetrics.MAX_TOTAL_FILE_SIZE))
.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();
Expand Down
Loading