Skip to content
Closed
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 @@ -70,6 +70,11 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca
<td>Gauge</td>
<td>Number of scanned manifest files in the last scan.</td>
</tr>
<tr>
<td>lastScanSkippedManifests</td>
<td>Gauge</td>
<td>Number of manifest files skipped by manifest level filtering in the last scan.</td>
</tr>
<tr>
<td>lastScanSkippedTableFiles</td>
<td>Gauge</td>
Expand All @@ -80,6 +85,16 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca
<td>Gauge</td>
<td>Resulted table files in the last scan.</td>
</tr>
<tr>
<td>lastScanResultedTableFilesSize</td>
<td>Gauge</td>
<td>Total size in bytes of the resulted table files to be read in the last scan.</td>
</tr>
<tr>
<td>lastScanResultedRecordCount</td>
<td>Gauge</td>
<td>Total number of records in the resulted table files to be read in the last scan.</td>
</tr>
</tbody>
</table>

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -319,11 +319,15 @@ public Plan plan() {
manifestsResult.allManifests.stream()
.mapToLong(f -> f.numAddedFiles() - f.numDeletedFiles())
.sum();
// The size and record count of what will be read are not folded here: a DELTA plan
// holds both ADD and DELETE entries, and whether the DELETE entries get read depends
// on the consumer. SnapshotReaderImpl reports them once it has made that choice.
scanMetrics.reportScan(
new ScanStats(
scanDuration,
snapshot == null ? 0 : snapshot.id(),
manifests.size(),
manifestsResult.allManifests.size() - manifests.size(),
allDataFiles - result.size(),
result.size()));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,12 @@ public class ScanMetrics {
public static final String SCAN_DURATION = "scanDuration";
public static final String LAST_SCANNED_SNAPSHOT_ID = "lastScannedSnapshotId";
public static final String LAST_SCANNED_MANIFESTS = "lastScannedManifests";
public static final String LAST_SCAN_SKIPPED_MANIFESTS = "lastScanSkippedManifests";
public static final String LAST_SCAN_SKIPPED_TABLE_FILES = "lastScanSkippedTableFiles";
public static final String LAST_SCAN_RESULTED_TABLE_FILES = "lastScanResultedTableFiles";
public static final String LAST_SCAN_RESULTED_TABLE_FILES_SIZE =
"lastScanResultedTableFilesSize";
public static final String LAST_SCAN_RESULTED_RECORD_COUNT = "lastScanResultedRecordCount";
public static final String MANIFEST_HIT_CACHE = "manifestHitCache";
public static final String MANIFEST_MISSED_CACHE = "manifestMissedCache";
public static final String DVMETA_HIT_CACHE = "dvMetaHitCache";
Expand All @@ -46,6 +50,12 @@ public class ScanMetrics {

private ScanStats latestScan;

// Reported separately from the scan itself: which entries of a plan are read depends on the
// consumer (a normal read takes the ADD entries, a change read takes DELETE entries as well),
// so the size and record count are folded after that choice is made, not at plan time.
private long latestResultedTableFilesSize;
private long latestResultedRecordCount;

public ScanMetrics(MetricRegistry registry, String tableName) {
metricGroup = registry.createTableMetricGroup(GROUP_NAME, tableName);
metricGroup.gauge(
Expand All @@ -59,12 +69,17 @@ public ScanMetrics(MetricRegistry registry, String tableName) {
metricGroup.gauge(
LAST_SCANNED_MANIFESTS,
() -> latestScan == null ? 0L : latestScan.getScannedManifests());
metricGroup.gauge(
LAST_SCAN_SKIPPED_MANIFESTS,
() -> latestScan == null ? 0L : latestScan.getSkippedManifests());
metricGroup.gauge(
LAST_SCAN_SKIPPED_TABLE_FILES,
() -> latestScan == null ? 0L : latestScan.getSkippedTableFiles());
metricGroup.gauge(
LAST_SCAN_RESULTED_TABLE_FILES,
() -> latestScan == null ? 0L : latestScan.getResultedTableFiles());
metricGroup.gauge(LAST_SCAN_RESULTED_TABLE_FILES_SIZE, () -> latestResultedTableFilesSize);
metricGroup.gauge(LAST_SCAN_RESULTED_RECORD_COUNT, () -> latestResultedRecordCount);
metricGroup.gauge(MANIFEST_HIT_CACHE, () -> cacheMetrics.getHitObject().get());
metricGroup.gauge(MANIFEST_MISSED_CACHE, () -> cacheMetrics.getMissedObject().get());
metricGroup.gauge(DVMETA_HIT_CACHE, () -> dvMetaCacheMetrics.getHitObject().get());
Expand All @@ -81,6 +96,15 @@ public void reportScan(ScanStats scanStats) {
durationHistogram.update(scanStats.getDuration());
}

/**
* Reports the total size and record count of the data files the consumer of the latest plan
* will actually read. Called after the reader has decided which entries it takes from the plan.
*/
public void reportResultedFiles(long tableFilesSize, long recordCount) {
latestResultedTableFilesSize = tableFilesSize;
latestResultedRecordCount = recordCount;
}

public CacheMetrics getCacheMetrics() {
return cacheMetrics;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,18 +26,21 @@ public class ScanStats {
private final long duration;
private final long scannedSnapshotId;
private final long scannedManifests;
private final long skippedManifests;
private final long skippedTableFiles;
private final long resultedTableFiles;

public ScanStats(
long duration,
long scannedSnapshotId,
long scannedManifests,
long skippedManifests,
long skippedTableFiles,
long resultedTableFiles) {
this.duration = duration;
this.scannedSnapshotId = scannedSnapshotId;
this.scannedManifests = scannedManifests;
this.skippedManifests = skippedManifests;
this.skippedTableFiles = skippedTableFiles;
this.resultedTableFiles = resultedTableFiles;
}
Expand All @@ -52,6 +55,11 @@ protected long getScannedManifests() {
return scannedManifests;
}

@VisibleForTesting
protected long getSkippedManifests() {
return skippedManifests;
}

@VisibleForTesting
protected long getSkippedTableFiles() {
return skippedTableFiles;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ public class SnapshotReaderImpl implements SnapshotReader {
private boolean hasNonPartitionFilter;
private RecordComparator lazyPartitionComparator;
private CacheMetrics dvMetaCacheMetrics;
@Nullable private ScanMetrics scanMetrics;

public SnapshotReaderImpl(
FileStoreScan scan,
Expand Down Expand Up @@ -322,12 +323,34 @@ public SnapshotReader withBucketFilter(Filter<Integer> bucketFilter) {

@Override
public SnapshotReader withMetricRegistry(MetricRegistry registry) {
ScanMetrics scanMetrics = new ScanMetrics(registry, tableName);
scanMetrics = new ScanMetrics(registry, tableName);
dvMetaCacheMetrics = scanMetrics.getDvMetaCacheMetrics();
scan.withMetrics(scanMetrics);
return this;
}

/**
* Reports the size and record count of the data files this reader is about to hand out. Which
* entries of a plan are read is decided here, not in the scan: a normal read takes the ADD
* entries, a change read also reads the DELETE entries as its before files. Reporting from the
* reader keeps the metrics tied to what actually gets read.
*/
@SafeVarargs
private final void reportResultedFiles(List<ManifestEntry>... readEntries) {
if (scanMetrics == null) {
return;
}
long tableFilesSize = 0L;
long recordCount = 0L;
for (List<ManifestEntry> entries : readEntries) {
for (ManifestEntry entry : entries) {
tableFilesSize += entry.file().fileSize();
recordCount += entry.file().rowCount();
}
}
scanMetrics.reportResultedFiles(tableFilesSize, recordCount);
}

@Override
public SnapshotReader withRowRanges(List<Range> sortedPushdownRowRanges) {
scan.withRowRanges(sortedPushdownRowRanges);
Expand Down Expand Up @@ -394,8 +417,10 @@ public Plan read() {
FileStoreScan.Plan plan = scan.plan();
@Nullable Snapshot snapshot = plan.snapshot();

Map<BinaryRow, Map<Integer, List<ManifestEntry>>> grouped =
groupByPartFiles(plan.files(FileKind.ADD));
// a normal read only takes the ADD entries, even from a DELTA plan
List<ManifestEntry> addFiles = plan.files(FileKind.ADD);
reportResultedFiles(addFiles);
Map<BinaryRow, Map<Integer, List<ManifestEntry>>> grouped = groupByPartFiles(addFiles);
if (options.scanPlanSortPartition()) {
Map<BinaryRow, Map<Integer, List<ManifestEntry>>> sorted = new LinkedHashMap<>();
grouped.entrySet().stream()
Expand Down Expand Up @@ -492,10 +517,15 @@ public Plan readChanges() {
withMode(ScanMode.DELTA);
FileStoreScan.Plan plan = scan.plan();

List<ManifestEntry> beforeEntries = plan.files(FileKind.DELETE);
List<ManifestEntry> afterEntries = plan.files(FileKind.ADD);
// both sides are read: the DELETE entries are the before files of the change
reportResultedFiles(beforeEntries, afterEntries);

Map<BinaryRow, Map<Integer, List<ManifestEntry>>> beforeFiles =
groupByPartFiles(plan.files(FileKind.DELETE));
groupByPartFiles(beforeEntries);
Map<BinaryRow, Map<Integer, List<ManifestEntry>>> afterFiles =
groupByPartFiles(plan.files(FileKind.ADD));
groupByPartFiles(afterEntries);
LazyField<Snapshot> beforeSnapshot =
new LazyField<>(() -> snapshotManager.snapshot(plan.snapshot().id() - 1));
return toIncrementalPlan(
Expand Down Expand Up @@ -604,10 +634,15 @@ private Plan toIncrementalPlan(
public Plan readIncrementalDiff(Snapshot before) {
withMode(ScanMode.ALL);
FileStoreScan.Plan plan = scan.plan();
List<ManifestEntry> afterEntries = plan.files(FileKind.ADD);
List<ManifestEntry> beforeEntries = scan.withSnapshot(before).plan().files(FileKind.ADD);
// two independent scans; the ADD entries of both snapshots are read to compute the diff
reportResultedFiles(beforeEntries, afterEntries);

Map<BinaryRow, Map<Integer, List<ManifestEntry>>> afterFiles =
groupByPartFiles(plan.files(FileKind.ADD));
groupByPartFiles(afterEntries);
Map<BinaryRow, Map<Integer, List<ManifestEntry>>> beforeFiles =
groupByPartFiles(scan.withSnapshot(before).plan().files(FileKind.ADD));
groupByPartFiles(beforeEntries);
TimeTravelUtil.checkRescaleBucketForIncrementalDiffQuery(
tableSchema, before, beforeFiles, plan.snapshot(), afterFiles);
return toIncrementalPlan(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,11 @@ public void testGenericMetricsRegistration() {
ScanMetrics.SCAN_DURATION,
ScanMetrics.LAST_SCANNED_SNAPSHOT_ID,
ScanMetrics.LAST_SCANNED_MANIFESTS,
ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS,
ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES,
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES,
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE,
ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT,
ScanMetrics.MANIFEST_HIT_CACHE,
ScanMetrics.MANIFEST_MISSED_CACHE,
ScanMetrics.DVMETA_HIT_CACHE,
Expand All @@ -71,20 +74,32 @@ public void testMetricsAreUpdated() {
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_SNAPSHOT_ID);
Gauge<Long> lastScannedManifests =
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_MANIFESTS);
Gauge<Long> lastScanSkippedManifests =
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS);
Gauge<Long> lastScanSkippedTableFiles =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES);
Gauge<Long> lastScanResultedTableFiles =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES);
Gauge<Long> lastScanResultedTableFilesSize =
(Gauge<Long>)
registeredGenericMetrics.get(
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE);
Gauge<Long> lastScanResultedRecordCount =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT);

assertThat(lastScanDuration.getValue()).isEqualTo(0);
assertThat(lastScannedSnapshotId.getValue()).isEqualTo(0);
assertThat(scanDuration.getCount()).isEqualTo(0);
assertThat(scanDuration.getStatistics().size()).isEqualTo(0);
assertThat(lastScannedManifests.getValue()).isEqualTo(0);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(0);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(0);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(0);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(0);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(0);

// report once
reportOnce(scanMetrics);
Expand All @@ -101,8 +116,11 @@ public void testMetricsAreUpdated() {
assertThat(scanDuration.getStatistics().getMax()).isEqualTo(200);
assertThat(scanDuration.getStatistics().getStdDev()).isEqualTo(0);
assertThat(lastScannedManifests.getValue()).isEqualTo(20);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(5);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(25);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(10);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(1024);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(100);

// report again
reportAgain(scanMetrics);
Expand All @@ -119,18 +137,21 @@ public void testMetricsAreUpdated() {
assertThat(scanDuration.getStatistics().getMax()).isEqualTo(500);
assertThat(scanDuration.getStatistics().getStdDev()).isCloseTo(212.132, offset(0.001));
assertThat(lastScannedManifests.getValue()).isEqualTo(22);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(7);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(30);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(8);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(2048);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(200);
}

private void reportOnce(ScanMetrics scanMetrics) {
ScanStats scanStats = new ScanStats(200, 1L, 20, 25, 10);
scanMetrics.reportScan(scanStats);
scanMetrics.reportScan(new ScanStats(200, 1L, 20, 5, 25, 10));
scanMetrics.reportResultedFiles(1024, 100);
}

private void reportAgain(ScanMetrics scanMetrics) {
ScanStats scanStats = new ScanStats(500, 2L, 22, 30, 8);
scanMetrics.reportScan(scanStats);
scanMetrics.reportScan(new ScanStats(500, 2L, 22, 7, 30, 8));
scanMetrics.reportResultedFiles(2048, 200);
}

private ScanMetrics getScanMetrics() {
Expand Down
Loading
Loading